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:
- 在
ConsumerImpl的__init__方法中调用了threadInit(),该方法会创建并启动一个新线程(Thread-1)执行_consumerRunner的无限循环; - 在
QueueDemo.py的主函数末尾,又直接调用consumer1._consumerRunner(),这会让主线程(MainThread)也进入同一个无限循环。
两个线程同时拉取消息、更新偏移量,自然会导致同一条消息被多次消费。
解决方法(二选一)
- 方案一:删除
QueueDemo.py末尾的consumer1._consumerRunner()调用,只保留threadInit()中启动的后台线程; - 方案二:去掉
ConsumerImpl.__init__里的self.threadInit()调用,只在你指定的单线程中手动调用_consumerRunner。
额外注意
如果后续需要支持多线程消费场景,你当前的偏移量更新操作__setTopicOffset没有线程安全保护,多个线程同时读写__topicVsOffset会出现数据不一致问题,需要用threading.Lock来包裹共享资源的读写逻辑。
内容的提问来源于stack exchange,提问作者Pygirl
相关产品推荐
相关产品推荐

