MongoDB变更流异步任务非等待启动与并发处理异常排查
解决MongoDB变更流异步并发处理问题
问题根源分析
- 直接创建任务不保存也不等待:
asyncio.create_task创建的任务如果没有被引用,可能被垃圾回收器回收;另外如果用同步for遍历变更流,会阻塞事件循环,任务根本得不到调度执行的机会。 - 原地await任务:会强制串行执行,完全失去并发能力。
- 收集任务但循环不退出:变更流是无限流,
for循环永远不会结束,后续的统一await逻辑永远执行不到,而且代码中还犯了一个错误——创建的任务没有添加到tasks列表中,就算循环结束也无任务可等待。
正确的并发处理实现
使用异步迭代变更流,维护任务列表并定期清理已完成任务,同时可限制并发数避免资源耗尽:
import asyncio async def handle_collection_changes( *, change_stream, handler, handler_args, db, collection_name, service_name, ): tasks = [] max_concurrent_tasks = 10 # 根据实际情况调整并发上限 try: # 异步遍历变更流,避免阻塞事件循环 async for change in change_stream: # 创建异步任务并加入列表 task = asyncio.create_task(handler(change, *handler_args)) tasks.append(task) # 保存恢复令牌 save_resume_token(db, collection_name, change["_id"], service_name) # 清理已完成的任务,释放资源 tasks = [t for t in tasks if not t.done()] # 达到并发上限时,等待至少一个任务完成再继续 if len(tasks) >= max_concurrent_tasks: done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED) tasks = list(pending) except KeyboardInterrupt: log_debug("keyboard interrupt detected, closing stream") except Exception as e: log_critical(f"unexpected error in change stream: {repr(e)}") finally: # 程序退出前,等待所有未完成的任务执行完毕 if tasks: await asyncio.gather(*tasks, return_exceptions=True)
关键要点说明
- 异步迭代:必须用
async for遍历MongoDB异步变更流(需配合异步驱动如Motor),否则同步遍历会完全阻塞事件循环,异步任务无法调度。 - 任务生命周期管理:维护任务列表避免GC回收,定期清理已完成任务防止内存泄漏。
- 并发数控制:通过
max_concurrent_tasks限制同时运行的任务数,避免系统资源被耗尽。 - 优雅退出:在
finally块中等待剩余任务完成,确保所有已接收的变更都被处理完毕。
内容的提问来源于stack exchange,提问作者SaratAngajala
相关产品推荐
相关产品推荐

