Confluent Kafka Python consumer.poll()执行后退出代码求助
Kafka消费者
poll()执行后直接退出问题排查与解决 问题现象
将Kafka视频流处理管道改为图片输入后,消费者执行consumer.poll()时代码直接退出,未执行后续逻辑。控制台仅打印entered,无poll()返回值输出。
Kafka容器日志
kafka1 | [2023-01-10 07:21:30,719] INFO [GroupCoordinator 1]: Dynamic member with unknown member id joins group kafka-cam in Empty state. Created a new member id rdkafka-f0ea3fd0-2eef-494d-bbe5-1acd9db067c4 and request the member to rejoin with this id. (kafka.coordinator.group.GroupCoordinator) kafka1 | [2023-01-10 07:21:30,721] INFO [GroupCoordinator 1]: Preparing to rebalance group kafka-cam in state PreparingRebalance with old generation 0 (__consumer_offsets-24) (reason: Adding new member rdkafka-f0ea3fd0-2eef-494d-bbe5-1acd9db067c4 with group instance id None; client reason: not provided) (kafka.coordinator.group.GroupCoordinator) kafka1 | [2023-01-10 07:21:33,722] INFO [GroupCoordinator 1]: Stabilized group kafka-cam generation 1 (__consumer_offsets-24) with 1 members (kafka.coordinator.group.GroupCoordinator)(kafka.coordinator.group.GroupCoordinator)
消费者核心代码与执行输出
消费者循环逻辑:
while True: print("entered") msg = consumer.poll(100) print(msg) if msg == None: continue
执行consumer.py输出:
python consumer.py Creating consumer thread... Starting the consumer thread... entered
问题原因分析
- 守护线程被强制终止:消费者线程被设置为
daemon=True,主线程启动线程后无阻塞逻辑直接退出,导致守护线程还未完成poll()就被系统终止。 - Producer配置冗余错误:Producer配置中包含
group.id、enable.auto.commit等消费者专属配置,可能引发组协调异常。 - 消费者逻辑分支错误:错误处理分支顺序混乱,且未正确使用类成员变量(如直接访问全局
db而非self.db),可能触发未捕获异常导致线程退出。
解决方案
1. 修复主线程阻塞逻辑
修改消费者启动代码,让主线程等待消费者线程结束:
if __name__ == "__main__": load_dotenv() topic = ["multi-cam-stream"] MODEL_API = os.environ["MODEL_API"] MONGODB_URI = os.environ["MONGODB_URI"] client = MongoClient(MONGODB_URI) db = client.smart_agriculture print("Creating consumer thread...") consumer_thread = ConsumerThread(consumer_config, topic, 4, MODEL_API, db) print("Starting the consumer thread...") # 创建非守护线程并等待其结束 threads = [] t = threading.Thread(target=consumer_thread.read_data) t.daemon = False threads.append(t) t.start() # 主线程等待线程完成 for t in threads: t.join()
2. 修正Producer配置
移除Producer中不属于生产者的配置项:
# producer_config.py config = { 'bootstrap.servers': '127.0.0.1:9092', }
3. 修复消费者逻辑错误
调整错误处理分支顺序,正确使用类成员变量,并清理批次处理后的缓存:
def run(self, consumer, msg_count): try: coll = self.db.outputs # 使用类成员变量self.db imgs = [] while True: print("entered") msg = consumer.poll(timeout=3) print(msg) if msg is None: print("msg is empty") continue # 统一处理消息错误 if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: print(f'End of partition reached {msg.topic()}/{msg.partition()}') else: raise KafkaException(msg.error()) else: # 处理正常消息 img_bytes = msg.value() imgs.append(img_bytes) print(len(imgs)) msg_count += 1 if msg_count % self.batch_size == 0: print(len(imgs)) files = { f"{_id}": img for _id, img in enumerate(imgs) } # 调用模型API response = requests.post(self.model_api, files=files) results = json.loads(response.content) for _id in files: img_str = base64.b64encode(files[_id]) img_doc = { "image": img_str, "result": results[int(_id)] } coll.insert_one(img_doc) consumer.commit(asynchronous=False) imgs.clear() # 清空批次缓存 print(msg_count) except KeyboardInterrupt: print("Detected Keyboard Interrupt. Quitting...") finally: print("closing") consumer.close()
完整代码参考
Producer文件(producer.py)
from producer_config import config as producer_config from confluent_kafka import Producer import os import concurrent.futures import cv2 from utils import * from glob import glob class ProducerThread: def __init__(self, config): self.producer = Producer(config) def publishImage(self, img_path): img = cv2.imread(img_path) print(img.shape) img_name = os.path.basename(img_path).split(".")[0] img_bytes = serializeImg(img) self.producer.produce( topic = "multi-cam-stream", value = img_bytes, on_delivery = delivery_report, timestamp = 0, headers = { "img_name": str.encode(img_name) } ) self.producer.poll(0.5) return def start(self, img_paths): with concurrent.futures.ThreadPoolExecutor() as executor: executor.map(self.publishImage, img_paths) self.producer.flush() # 推送队列中剩余消息 print("Finished") if __name__ == "__main__": img_dir = "imgs/" img_paths = glob(img_dir + "*.JPG") producer_thread = ProducerThread(producer_config) producer_thread.start(img_paths)
Producer配置(producer_config.py)
config = { 'bootstrap.servers': '127.0.0.1:9092', }
Consumer文件(consumer.py)
import threading from confluent_kafka import Consumer, KafkaError, KafkaException from consumer_config import config as consumer_config from utils import * from pymongo import MongoClient from dotenv import load_dotenv import os import base64 import requests import json class ConsumerThread: def __init__(self, config, topic, batch_size, model_api, db): self.config = config self.topic = topic self.batch_size = batch_size self.model_api = model_api self.db = db def read_data(self): consumer = Consumer(self.config) consumer.subscribe(self.topic) self.run(consumer,0) def run(self, consumer, msg_count): try: coll = self.db.outputs imgs = [] while True: print("entered") msg = consumer.poll(timeout=3) print(msg) if msg is None: print("msg is empty") continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: print(f'End of partition reached {msg.topic()}/{msg.partition()}') else: raise KafkaException(msg.error()) else: img_bytes = msg.value() imgs.append(img_bytes) print(len(imgs)) msg_count += 1 if msg_count % self.batch_size == 0: print(len(imgs)) files = { f"{_id}": img for _id, img in enumerate(imgs) } response = requests.post(self.model_api, files=files) results = json.loads(response.content) for _id in files: img_str = base64.b64encode(files[_id]) img_doc = { "image": img_str, "result": results[int(_id)] } coll.insert_one(img_doc) consumer.commit(asynchronous=False) imgs.clear() print(msg_count) except KeyboardInterrupt: print("Detected Keyboard Interrupt. Quitting...") finally: print("closing") consumer.close() def start(self, numThreads): threads = [] for _ in range(numThreads): t = threading.Thread(target=self.read_data) t.daemon = False threads.append(t) t.start() return threads if __name__ == "__main__": load_dotenv() topic = ["multi-cam-stream"] MODEL_API = os.environ["MODEL_API"] MONGODB_URI = os.environ["MONGODB_URI"] client = MongoClient(MONGODB_URI) db = client.smart_agriculture print("Creating consumer thread...") consumer_thread = ConsumerThread(consumer_config, topic, 4, MODEL_API, db) print("Starting the consumer thread...") threads = consumer_thread.start(1) # 主线程等待所有线程结束 for t in threads: t.join()
Consumer配置(consumer_config.py)
config = { 'bootstrap.servers': '127.0.0.1:9092', 'group.id': 'kafka-cam', 'enable.auto.commit': False, 'auto.offset.reset': 'earliest', # 替换旧的default.topic.config格式 'max.poll.interval.ms': 20000, 'session.timeout.ms': 10000, 'fetch.message.max.bytes': 10000000, 'max.partition.fetch.bytes': 1000000 }
Utils.py
import logging import cv2 import json logging.basicConfig(level=logging.INFO, format='%(name)s - %(levelname)s - %(message)s') def delivery_report(err, msg): if err: logging.error("Failed to deliver message: {0}: {1}" .format(msg.value(), err.str())) else: logging.info(f"msg produced. \n" f"Topic: {msg.topic()} \n" + f"Partition: {msg.partition()} \n" + f"Offset: {msg.offset()} \n" + f"Timestamp: {msg.timestamp()} \n") def serializeImg(img): _, img_buffer_arr = cv2.imencode(".jpg", img) img_bytes = img_buffer_arr.tobytes() return img_bytes def jsonify(img): return json.dumps(img)
内容的提问来源于stack exchange,提问作者Eric
相关产品推荐
相关产品推荐

