C#中并行执行任务时等待指定任务完成的实现方案
解决方案:基于队列字典的Topic任务调度
你的思路完全可行,用字典维护不同Topic的任务队列,刚好能匹配「Topic A串行、Topic B并行」的差异化需求。下面是具体的落地实现思路和代码示例:
核心设计逻辑
用字典存储各Topic对应的任务调度单元:
- Topic A:绑定「任务队列+单消费线程」,保证同一Topic的任务按消息到达顺序串行执行,前一个任务完成后才启动下一个,新消息会自动加入队列等待
- Topic B:直接提交到线程池并行处理,无需额外队列限制,多个Topic B的任务可同时运行
代码实现示例(Python)
以常用的paho-mqtt客户端和线程池为例:
import paho.mqtt.client as mqtt from concurrent.futures import ThreadPoolExecutor import queue import threading # 初始化线程池,可根据服务器性能调整并发数 executor = ThreadPoolExecutor(max_workers=10) # 维护Topic对应的队列与锁:key为Topic,value是(队列对象, 线程锁) topic_task_config = { "Topic/A": (queue.Queue(), threading.Lock()), # Topic B无需专用队列,直接走线程池并行 } def time_consuming_task(msg_content): """模拟耗时1秒的业务处理方法""" import time print(f"启动任务: {msg_content}") time.sleep(1) print(f"完成任务: {msg_content}") def consume_topic_a_queue(): """Topic A的专属消费线程,保证任务串行执行""" q, lock = topic_task_config["Topic/A"] while True: # 阻塞等待队列中的新任务 msg = q.get() try: time_consuming_task(msg) finally: # 标记当前任务处理完成,队列可取出下一个任务 q.task_done() def on_mqtt_message(client, userdata, msg): topic = msg.topic payload = msg.payload.decode("utf-8") if topic == "Topic/A": q, lock = topic_task_config["Topic/A"] # 加锁保证多线程下队列操作的安全性 with lock: q.put(payload) elif topic == "Topic/B": # 直接提交到线程池,并行执行 executor.submit(time_consuming_task, payload) # 初始化MQTT客户端 client = mqtt.Client() client.on_message = on_mqtt_message # 启动Topic A的串行消费线程(后台守护线程) threading.Thread(target=consume_topic_a_queue, daemon=True).start() # 连接MQTT服务器并订阅目标Topic client.connect("你的MQTT Broker地址", 1883, 60) client.subscribe("Topic/A") client.subscribe("Topic/B") client.loop_forever()
关键细节说明
- Topic A串行保障:单独启动的消费线程会持续从队列取任务,只有当前任务执行完毕,才会处理下一个,完全符合「上一个任务完成再启动新任务」的要求,新消息会被自动加入队列等待,不会丢失
- Topic B并行逻辑:直接将任务提交到线程池,线程池会自动分配空闲线程处理,多个Topic B任务可同时运行,互不干扰
- 线程安全:Topic A的队列操作加了线程锁,避免MQTT回调线程(可能存在多个)同时写入队列导致的异常
- 扩展性:后续新增串行Topic时,只需在
topic_task_config中添加对应的队列和锁,再启动一个消费线程即可;新增并行Topic则直接复用现有逻辑
可选优化方向
- 可以给Topic B也添加任务队列,同时启动多个消费线程,既保证消息不丢失,又能控制并行数,避免线程池资源耗尽
- 给队列设置最大长度,防止消息堆积过多导致内存溢出,超出长度时可选择丢弃最老消息或拒绝新消息(根据业务需求调整)
内容的提问来源于stack exchange,提问作者Max
相关产品推荐
相关产品推荐

