Update
This commit is contained in:
parent
c5b6582c64
commit
c59d406493
@ -78,6 +78,8 @@ class KafkaConsumerThread(threading.Thread):
|
||||
self.shutdown = threading.Event()
|
||||
|
||||
def run(self):
|
||||
consumer = None
|
||||
try:
|
||||
consumer = KafkaConsumer(
|
||||
self.topic,
|
||||
bootstrap_servers=self.bootstrap_servers,
|
||||
@ -85,6 +87,8 @@ class KafkaConsumerThread(threading.Thread):
|
||||
key_deserializer=lambda x: decode_key(x),
|
||||
value_deserializer=lambda x: decode_value(x)
|
||||
)
|
||||
except:
|
||||
self.stop()
|
||||
|
||||
logger.info("Kafka Consumer started")
|
||||
|
||||
|
||||
Loading…
Reference in New Issue
Block a user