Skip to content

Commit b94e692

Browse files
authored
Merge pull request #23 from esakik/development
RELEASE 2023-12-30
2 parents 23df21f + 588f40e commit b94e692

1 file changed

Lines changed: 16 additions & 9 deletions

File tree

  • beam_mysql/connector

beam_mysql/connector/io.py

Lines changed: 16 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -118,7 +118,7 @@ def __init__(
118118

119119
def start_bundle(self):
120120
self._build_value()
121-
self._queries = []
121+
self._columns_and_values = dict()
122122

123123
def process(self, element: Dict, *args, **kwargs):
124124
columns = []
@@ -137,18 +137,25 @@ def process(self, element: Dict, *args, **kwargs):
137137
]
138138
)
139139

140-
query = f"INSERT INTO {self._config['database']}.{self._table}({column_str}) VALUES({value_str});"
140+
if column_str not in self._columns_and_values:
141+
self._columns_and_values[column_str] = []
141142

142-
self._queries.append(query)
143+
self._columns_and_values[column_str].append(f"({value_str})")
143144

144-
if len(self._queries) > self._batch_size:
145-
self._client.record_loader("\n".join(self._queries))
146-
self._queries.clear()
145+
if len(self._columns_and_values) > self._batch_size:
146+
for column_str in self._columns_and_values.keys():
147+
if len(self._columns_and_values[column_str]) > 0:
148+
self._client.record_loader(self._build_query(column_str, self._columns_and_values[column_str]))
149+
self._columns_and_values[column_str].clear()
147150

148151
def finish_bundle(self):
149-
if len(self._queries):
150-
self._client.record_loader("\n".join(self._queries))
151-
self._queries.clear()
152+
for column_str in self._columns_and_values.keys():
153+
if len(self._columns_and_values[column_str]) > 0:
154+
self._client.record_loader(self._build_query(column_str, self._columns_and_values[column_str]))
155+
self._columns_and_values[column_str].clear()
156+
157+
def _build_query(self, column_str, values_str):
158+
return f"INSERT INTO {self._config['database']}.{self._table}({column_str}) VALUES {','.join(values_str)};"
152159

153160
def _build_value(self):
154161
for k, v in self._config.items():

0 commit comments

Comments
 (0)