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

如何搭建异步监听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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 08:57:41