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

Python类Kafka内存队列:如何确保consumerRunner单线程运行避免重复消费?

问题描述

我正在用Python实现类Kafka的内存队列,用到了可重入锁和线程相关概念,刚接触Python线程。已经实现了consumer组件,能订阅topic并读取消息,整体功能正常,但有线程相关的疑问:

我期望consumerRunner以单线程运行,但程序输出显示MainThread和Thread-1两个线程都在执行这个方法,导致message2和message4被重复消费。如果是单线程运行,消息应该只被消费一次,我知道主线程是默认线程,但这种重复打印正常吗?

程序输出:

(py3_8) ninjakx@Kritis-MacBook-Pro kafka % python queueDemo.py
Msg: message1 has been published to topic: topic1
Msg: message2 has been published to topic: topic1
Msg: message3 has been published to topic: topic2
Msg: message4 has been published to topic: topic1
Msg: message1 has been consumed by consumer: consumer1 at offset: 0 with current thread: MainThread
Msg: message3 has been consumed by consumer: consumer1 at offset: 0 with current thread: MainThread
Msg: message2 has been consumed by consumer: consumer1 at offset: 1 with current thread: Thread-1
Msg: message2 has been consumed by consumer: consumer1 at offset: 1 with current thread: MainThread
Msg: message4 has been consumed by consumer: consumer1 at offset: 2 with current thread: Thread-1
Msg: message4 has been consumed by consumer: consumer1 at offset: 2 with current thread: MainThread

ConsumerImpl.py

import zope.interface
from ..interface.iConsumer import iConsumer
from collections import OrderedDict
from mediator.QueueMediatorImpl import QueueMediatorImpl

import threading 
from threading import Thread
import time
 
@zope.interface.implementer(iConsumer)
class ConsumerImpl:
    # will keep all the topic it has subscribed to and their offset

    def __init__(self, consumerName:str):
        self.__consumerName = consumerName
        self.__topicList = []
        self.__topicVsOffset = OrderedDict()
        self.__queueMediator = QueueMediatorImpl()
        self.threadInit()


    def threadInit(self):
        thread = Thread(target = self._consumerRunner)
        thread.start()
        # thread.join()
        # print("thread finished...exiting")

    def __getConsumerName(self):
        return self.__consumerName

    def __getQueueMediator(self):
        return self.__queueMediator

    def __getSubscribedTopics(self)->list:
        return self.__topicList

    def __setTopicOffset(self, topicName:str, offset:int)->int:
        self.__topicVsOffset[topicName] = offset

    def __getTopicOffset(self, topicName:str)->int:
        return self.__topicVsOffset[topicName]

    def __addToTopicList(self, topicName:str)->None:
        self.__topicList.append(topicName)

    def _subToTopic(self, topicName:str):
        self.__addToTopicList(topicName)
        self.__topicVsOffset[topicName] = 0

    def __consumeMsg(self, msg:str, offset:int):
        print(f"Msg: {msg} has been consumed by consumer: {self.__getConsumerName()} at offset: {offset} with current thread: {threading.current_thread().name}")


    # pull based mechanism
    # running on single thread
    def _consumerRunner(self):
        while(True):
            for topicName in self.__getSubscribedTopics():
                curOffset = self.__getTopicOffset(topicName)
                qmd = self.__getQueueMediator()
                msg = qmd._readMsgIfPresent(topicName, curOffset)
                if msg is not None:
                    self.__consumeMsg(msg._getMessage(), curOffset)
                    curOffset += 1
                # update offset
                self.__setTopicOffset(topicName, curOffset)

            try:
                #sleep for 100 milliseconds
                #thread sleep
                # "sleep() makes the calling thread sleep until seconds seconds have elapsed or a signal arrives which is not ignored."
                time.sleep(0.1)
            except Exception as e:
                print(f"Error: {e}")

QueueDemo.py

from service.QueueServiceImpl import QueueServiceImpl

if __name__ == "__main__":
    queueService = QueueServiceImpl()
    producer1 = queueService._createProducer("producer1")
    producer2 = queueService._createProducer("producer2")
    producer3 = queueService._createProducer("producer3")
    producer4 = queueService._createProducer("producer4")

    consumer1 = queueService._createConsumer("consumer1")
    consumer2 = queueService._createConsumer("consumer2")
    consumer3 = queueService._createConsumer("consumer3")
    
    producer1._publishToTopic("topic1", "message1")
    producer1._publishToTopic("topic1", "message2")
    producer2._publishToTopic("topic2", "message3")
    producer1._publishToTopic("topic1", "message4")

    consumer1._subToTopic("topic1")
    consumer1._subToTopic("topic2")

    consumer1._consumerRunner()        
问题分析与解决

重复消费完全不正常,根源是你启动了两次_consumerRunner:

  1. 在ConsumerImpl的__init__方法中调用了threadInit(),该方法会创建并启动一个新线程(Thread-1)执行_consumerRunner的无限循环;
  2. 在QueueDemo.py的主函数末尾,又直接调用consumer1._consumerRunner(),这会让主线程(MainThread)也进入同一个无限循环。

两个线程同时拉取消息、更新偏移量,自然会导致同一条消息被多次消费。

解决方法(二选一)

  • 方案一:删除QueueDemo.py末尾的consumer1._consumerRunner()调用,只保留threadInit()中启动的后台线程;
  • 方案二:去掉ConsumerImpl.__init__里的self.threadInit()调用,只在你指定的单线程中手动调用_consumerRunner。

额外注意

如果后续需要支持多线程消费场景,你当前的偏移量更新操作__setTopicOffset没有线程安全保护,多个线程同时读写__topicVsOffset会出现数据不一致问题,需要用threading.Lock来包裹共享资源的读写逻辑。

内容的提问来源于stack exchange,提问作者Pygirl

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 11:54:51