如何搭建异步监听MongoDB Change Streams的Python生产级服务
异步监听MongoDB变更流的生产级实现方案
绝对懂你的困扰——官网和PyMongo文档里的同步轮询方案在生产环境确实拉胯,要么阻塞主线程,要么轮询间隔不好把控,很容易丢事件或者浪费资源。下面给你两个基于Python事件循环的异步实现方案,都是生产环境能直接用的:
方案一:用Motor(PyMongo官方异步驱动)实现原生异步监听
Motor是MongoDB官方推出的异步Python驱动,完美适配asyncio事件循环,是异步场景下的首选,性能和可靠性都更有保障。
实现示例
import asyncio import motor.motor_asyncio from bson.json_util import dumps import logging import time # 配置生产级日志,建议输出到文件或专业日志系统 logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s") logger = logging.getLogger(__name__) # 持久化resume token(生产环境建议存在MongoDB集合、Redis或配置中心,这里用文件做示例) def save_resume_token(token): if token: with open("resume_token.json", "w") as f: f.write(dumps(token)) def load_resume_token(): try: with open("resume_token.json", "r") as f: return dumps(f.read()) except FileNotFoundError: return None async def watch_collection(): client = motor.motor_asyncio.AsyncIOMotorClient("<YOUR-MONGO-CONNECT-STRING>") db = client.changestream collection = db.collection # 加载断点续传的token,避免服务重启后丢失事件 resume_token = load_resume_token() # 自定义过滤管道:只监听insert/update操作,减少无效处理 pipeline = [{"$match": {"operationType": {"$in": ["insert", "update"]}}}] while True: try: # 初始化变更流,支持断点续传 async with collection.watch(pipeline, resume_after=resume_token) as stream: async for change in stream: logger.info(f"捕获到变更: {dumps(change)}") # 这里替换成你的业务逻辑:比如发送到MQ、更新缓存等 # await process_change(change) # 实时保存resume token,进程崩溃后也能从断点恢复 save_resume_token(stream.resume_token) except Exception as e: logger.error(f"变更流监听异常: {str(e)},5秒后重试...") await asyncio.sleep(5) # 重试前重新加载最新的token,确保从正确断点开始 resume_token = load_resume_token() if __name__ == "__main__": loop = asyncio.get_event_loop() try: loop.run_until_complete(watch_collection()) except KeyboardInterrupt: logger.info("监听服务已手动停止")
生产环境核心注意点
- Resume Token持久化:这是生产环境必须做的,不然服务重启会从头监听,丢失中间的变更事件。
- 异常自动重连:无限循环+异常捕获是高可用的基础,网络波动、MongoDB节点切换都会导致连接断开,自动重连才能保证服务不中断。
- 过滤管道优化:通过
pipeline参数过滤掉不需要的操作类型或文档,减少不必要的计算和IO。
方案二:用asyncio线程池包装PyMongo同步监听(适合已用PyMongo的存量项目)
如果你的项目已经大量使用PyMongo,不想切换到Motor,可以用asyncio的线程池把同步的watch()包装成异步任务,避免阻塞主事件循环。
实现示例
import asyncio import pymongo from bson.json_util import dumps import logging from concurrent.futures import ThreadPoolExecutor import time logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s") logger = logging.getLogger(__name__) def save_resume_token(token): if token: with open("resume_token.json", "w") as f: f.write(dumps(token)) def load_resume_token(): try: with open("resume_token.json", "r") as f: return dumps(f.read()) except FileNotFoundError: return None def sync_watch_task(collection, resume_token): pipeline = [{"$match": {"operationType": {"$in": ["insert", "update"]}}}] while True: try: with collection.watch(pipeline, resume_after=resume_token) as stream: for change in stream: logger.info(f"捕获到变更: {dumps(change)}") # 同步处理业务逻辑,或者把事件放到异步队列里 # process_sync_change(change) save_resume_token(stream.resume_token) except Exception as e: logger.error(f"同步监听异常: {str(e)},5秒后重试...") time.sleep(5) resume_token = load_resume_token() async def async_watch_wrapper(): client = pymongo.MongoClient("<YOUR-MONGO-CONNECT-STRING>") collection = client.changestream.collection resume_token = load_resume_token() # 用单线程池跑同步监听任务,避免多线程导致事件乱序 executor = ThreadPoolExecutor(max_workers=1) loop = asyncio.get_event_loop() # 把同步任务扔到线程池,不阻塞主事件循环 await loop.run_in_executor(executor, sync_watch_task, collection, resume_token) if __name__ == "__main__": loop = asyncio.get_event_loop() try: loop.run_until_complete(async_watch_wrapper()) except KeyboardInterrupt: logger.info("监听服务已手动停止")
注意点
- 这个方案本质还是同步监听,只是把任务放到线程里隔离,性能上不如Motor原生异步,适合存量项目过渡。
- 线程池
max_workers必须设为1,因为变更流是有序的,多线程会导致事件处理乱序。
内容的提问来源于stack exchange,提问作者gustavz
相关产品推荐
相关产品推荐

