Update python
This commit is contained in:
parent
d866038826
commit
b81094f90e
@ -87,11 +87,13 @@ class KafkaConsumerThread(threading.Thread):
|
||||
key_deserializer=lambda x: decode_key(x),
|
||||
value_deserializer=lambda x: decode_value(x)
|
||||
)
|
||||
except:
|
||||
self.stop()
|
||||
return
|
||||
logger.info("Kafka Consumer started")
|
||||
|
||||
except:
|
||||
logger.exception("Kafka Consumer failed to start")
|
||||
self.stop()
|
||||
sys.exit(1)
|
||||
|
||||
logger.info("Kafka Consumer started")
|
||||
|
||||
while not self.shutdown.is_set():
|
||||
for message in consumer:
|
||||
|
||||
Loading…
Reference in New Issue
Block a user