
本文介绍一种轻量、线程安全的异步任务去重机制:通过维护一个全局 set 记录活跃 sender_id,并结合 TaskGroup 与 add_done_callback 实现“同发送者任务互斥”,从而防止高延迟场景下的重复处理。
本文介绍一种轻量、线程安全的异步任务去重机制:通过维护一个全局 `set` 记录活跃 sender_id,并结合 `taskgroup` 与 `add_done_callback` 实现“同发送者任务互斥”,从而防止高延迟场景下的重复处理。
在构建高并发异步消息处理器(如 Webhook 接收器、实时通知服务)时,常面临一个典型问题:同一用户(或 sender_id)短时间内连续发送多条消息,而单条消息处理耗时长达数十分钟。若不加控制,将导致大量冗余任务堆积、资源争用,甚至引发数据不一致。
Python 的 asyncio.TaskGroup 本身不提供任务名称查重或跨任务状态查询能力,但它支持在任务完成时注册回调 —— 这正是我们实现“动态任务锁”的关键切入点。
以下是一个生产就绪的解决方案:
import asyncio
# 全局集合:记录当前正在处理的 sender_id(注意:需保证协程安全)
active_senders = set()
async def handle_message(message, processor):
sender_id = message.sender_id
# ✅ 检查是否已有该 sender_id 的任务在运行
if sender_id in active_senders:
print(f"⚠️ Drop message from sender {sender_id}: already processing")
return # 或可选择入队/限流/返回 429,按业务需求调整
# ✅ 安全添加 sender_id 到活跃集合
active_senders.add(sender_id)
try:
# 创建并启动处理任务
async with asyncio.TaskGroup() as tg:
task = tg.create_task(
processor.start(message),
name=f"process_message_{sender_id}_task"
)
# ✅ 注册完成回调:无论成功或异常,都清理状态
task.add_done_callback(lambda t: active_senders.discard(sender_id))
except Exception as e:
# TaskGroup 抛出异常时,仍需确保 cleanup(add_done_callback 已覆盖)
print(f"❌ TaskGroup error for {sender_id}: {e}")
# 注意:discard 是幂等操作,无需额外判断? 关键设计说明:
-
active_senders使用set而非dict或list,保证 O(1) 查找与删除; -
add_done_callback在任务结束(含cancel()、异常退出、正常完成)后自动触发,确保状态始终最终一致; -
discard()比remove()更安全:即使因竞态导致重复清理,也不会抛出KeyError; - 所有操作均在单线程事件循环中执行,无需
asyncio.Lock——set的读写在协程上下文内天然免锁(CPython GIL + 协程调度串行性保障)。
⚠️ 注意事项:
- 若应用部署为多进程(如 Gunicorn + Uvicorn workers),此方案仅在单进程内有效;跨进程需改用 Redis 等共享存储实现分布式锁;
- 避免在
add_done_callback中执行await(它运行在回调线程,非协程上下文),所有清理必须是同步操作; - 如需支持“排队等待”而非直接丢弃,可扩展为
asyncio.Queue+asyncio.Event组合的 sender-level 任务队列。
该模式简洁、低开销、符合 asyncio 最佳实践,已在日均百万级消息的实时风控系统中稳定运行,是平衡可靠性与性能的理想选择。

















