改为手动提交

This commit is contained in:
wuaho 2021-10-26 10:40:44 +08:00
parent 4eea6929d7
commit 22fcd5e651
2 changed files with 6 additions and 2 deletions

2
app.py
View File

@ -79,7 +79,7 @@ class XProcess(Process):
else:
continue
transmitter.run()
transmitter.run(kafka_client)
while True:
time.sleep(5)

View File

@ -130,7 +130,7 @@ class Transmitter:
for key in del_keys:
del item[key]
def run(self):
def run(self, kafka_client):
for tb, buffer in self.check_send():
try:
data = [self.flat_data(x) for x in buffer.values()]
@ -141,3 +141,7 @@ class Transmitter:
except Exception as e:
self.log.error(traceback.format_exc())
buffer.clear()
try:
kafka_client.commit()
except Exception as e:
self.log.error(e)