我试图优雅地关闭一个kafka消费者,但是脚本阻止了心跳线程的停止。我如何才能优雅地关闭消费者与Kafkapython的sigterm。这就是我所做的
import logger as logging
import time
import sys
from kafka import KafkaConsumer
import numpy as np
import signal
log = logging.getLogger(__name__)
class Cons:
def __init__(self):
signal.signal(signal.SIGINT, self.sigterm_handler)
signal.signal(signal.SIGTERM, self.sigterm_handler)
self.consumer = KafkaConsumer('dummy-topic', group_id='poll-test', bootstrap_servers=['b1'])
def sigterm_handler(self, signum, frame):
log.info("Sigterm handler")
self.consumer.close(autocommit=False)
sys.exit(0)
def consume(self):
try:
while True:
records = self.consumer.poll(timeout_ms=500, max_records=500)
for topic_partition, consumer_records in records.items():
for record in consumer_records:
log.info("Got Record - {}".format(record))
#code to manually commit
except ValueError as e:
log.exception("exception")
if __name__ == '__main__':
c=Cons()
c.consume()
启用调试日志后,这就是我得到的输出,代码在此日志上被阻止。
^C2020-04-28 07:18:33,050 - MainThread - __main__ - INFO - Sigterm handler
2020-04-28 07:18:33,050 - MainThread - kafka.consumer.group - DEBUG - Closing the KafkaConsumer.
2020-04-28 07:18:33,051 - MainThread - kafka.coordinator - INFO - Stopping heartbeat thread
这背后的原因是什么?什么是关闭sigterm或sigint消费者的正确方法?
暂无答案!
目前还没有任何答案,快来回答吧!