You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

问题原因分析

  1. 守护线程被强制终止:消费者线程被设置为daemon=True,主线程启动线程后无阻塞逻辑直接退出,导致守护线程还未完成poll()就被系统终止。
  2. Producer配置冗余错误:Producer配置中包含group.id、enable.auto.commit等消费者专属配置,可能引发组协调异常。
  3. 消费者逻辑分支错误:错误处理分支顺序混乱,且未正确使用类成员变量(如直接访问全局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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 11:01:14