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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 05:50:25