ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

消息队列的协作模式

消息队列的协作模式 综述消息队列的协作模式分如下2种: (1) 只有一个主题的推拉模式。只有一个主题消息生产者将消息推送到(一个或多个)队列中消息消费者从队列中拉取消息进行处理。 (2) 发布/订阅模式。会有多个主题。消息被发送到一个主题多个消费者订阅该主题并从对应的队列中接收消息。 一、只有一个主题的推拉模式举例from abc import ABC import threading import queue import copy import time #生产者/消费者基类 class BaseComponent(ABC): def __init__(self): self._producers [] self._consumers [] def bind_producer(self, upstream): self._producers.append(upstream) def bind_consumer(self, downstream): self._consumers.append(downstream) class Send_Recv(BaseComponent): def __init__(self): super().__init__() self._work_flag None def start(self): self._work_flag True def send_msg(self, msg): pass def stop(self): self._work_flag False class SendConsumer(BaseComponent): def __init__(self): super().__init__() self._send_thread None self._run_flag None def start(self): self._run_flag True self._send_thread threading.Thread(nameMainThread, targetself._main_process) self._send_thread.start() def _main_process(self): while self._run_flag: data[0x02] received_data self._producers[0].send_msg(data) for consumer in self._consumers: consumer.consume(received_data) def consume(self, msg): #消费数据 pass def stop(self) - None: self._run_flag False class DataParseConsumer_1(BaseComponent): def __init__(self): super().__init__() self._data_queue queue.Queue() def consume(self, raw_data): #把数据放在缓存队列中 self._data_queue.put(copy.copy(raw_data)) pass def _get_data(self): while self._run_flag: try: data self._data_queue.get(blockFalse) except queue.Empty as e: time.sleep(0.001) continue # 处理data def start(self): self._run_flag True super().start() pass def stop(self): super().stop() class DataParseConsumer_2(BaseComponent): def __init__(self): super().__init__() self._data_queue queue.Queue() def consume(self, raw_data): #把数据放在缓存队列中 self._data_queue.put(copy.copy(raw_data)) pass def _get_data(self): while self._run_flag: try: data self._data_queue.get(blockFalse) except queue.Empty as e: time.sleep(0.001) continue # 处理data def start(self): self._run_flag True super().start() pass def stop(self): super().stop() def bind_producer_and_consumer(producer, consumer): producer.bind_consumer(consumer) consumer.bind_producer(producer) #生成生产者/消费者对象 data_parse_consume_1 DataParseConsumer_1() data_parse_consume_2 DataParseConsumer_2() p Send_Recv() p_c SendConsumer() #关联生产者/消费者对象 bind_producer_and_consumer(p, p_c) bind_producer_and_consumer(p_c, data_parse_consume_1) bind_producer_and_consumer(p_c, data_parse_consume_2)上述举例中对于p_c来说data_parse_consume_1和data_parse_consume_2都是消费者属于同一个主题。 p_c的_consumers有同属于一个主题的两个消费者即[data_parse_consume_1, data_parse_consume_2]。二、下面举例两个主题的消息队列import threading import time import queue from typing import Dict, List class MessageBroker: 发布订阅消息代理中心维护主题-订阅者队列映射 def __init__(self): # key: topic主题, value: 该主题下所有订阅者的队列列表 self._topics: Dict[str, List[queue.Queue]] {} self._lock threading.Lock() self._exit_sentinel None # 哨兵退出标记 def subscribe(self, topic: str) - queue.Queue: 订阅某个主题返回该订阅者专属队列 sub_q queue.Queue() with self._lock: if topic not in self._topics: self._topics[topic] [] self._topics[topic].append(sub_q) return sub_q def unsubscribe(self, topic: str, sub_q: queue.Queue): 取消订阅 with self._lock: if topic in self._topics and sub_q in self._topics[topic]: self._topics[topic].remove(sub_q) def publish(self, topic: str, msg): 发布消息广播给该主题全部订阅者每个订阅者队列放入一份消息副本 with self._lock: if topic not in self._topics: return for q in self._topics[topic]: q.put(msg) def broadcast_exit(self, topic: str): 给指定主题所有订阅者发送退出哨兵消息通知线程退出 self.publish(topic, self._exit_sentinel) class Publisher(threading.Thread): 发布者(生产者)面向对象封装向指定主题发布消息 def __init__(self, broker: MessageBroker, topic: str, publish_cnt: int, name: str): super().__init__() self.broker broker self.topic topic self.publish_cnt publish_cnt self.name name def run(self): for i in range(self.publish_cnt): msg {publisher: self.name, data: fdata_{i}} self.broker.publish(self.topic, msg) print(f【{self.name}】发布主题[{self.topic}] → {msg}) time.sleep(0.3) print(f {self.name} 发布完成 ) class Subscriber(threading.Thread): 订阅者(消费者)订阅主题接收广播消息 def __init__(self, broker: MessageBroker, topic: str, name: str): super().__init__() self.broker broker self.topic topic self.name name self.sub_queue broker.subscribe(topic) # 订阅拿到专属队列 def run(self): while True: msg self.sub_queue.get() # 读到哨兵退出循环 if msg is None: print(f {self.name} 收到退出信号结束订阅 ) self.sub_queue.task_done() break print(f【{self.name}】收到主题[{self.topic}]消息: {msg}) time.sleep(0.5) # 模拟消费耗时 self.sub_queue.task_done() def close(self): self.broker.unsubscribe(self.topic, self.sub_queue) if __name__ __main__: # 1. 创建消息代理中心 broker MessageBroker() topic_name sensor_topic # 2. 创建发布者生产者 pub1 Publisher(broker, topictopic_name, publish_cnt4, name发布者‑A) # 3. 创建2个订阅者订阅同一个主题发布一条消息两个订阅者都会收到 sub1 Subscriber(broker, topictopic_name, name订阅者‑1) sub2 Subscriber(broker, topictopic_name, name订阅者‑2) # 启动线程 pub1.start() sub1.start() sub2.start() # 等待发布者全部消息发布完毕 pub1.join() # 广播退出信号通知该主题所有订阅者退出 broker.broadcast_exit(topic_name) # 等待订阅者线程安全结束 sub1.join() sub2.join() sub1.close() sub2.close() print(\n 整个发布‑订阅程序全部结束)broker是消息代理中心其_topics可以有多个主题。_topics{主题1:[订阅者1的queue1, 订阅者1的queue2], 主题2:[订阅者3的queue3, 订阅者4的queue4]}
返回列表