Apache Airflow 集成 MongoDB 实战:MongoHook、MongoSensor 与 Airflow Connection 配置指南
发布时间:2026/10/9 2:20:28 锦皓数字建站

【免费下载链接】context-hub项目地址https://gitcode.com/gh_mirrors/co/context-hub点击查看免费下载apache-airflow-providers-mongo是 Apache Airflow 的官方 MongoDB Provider它把 Airflow 的 Connection 抽象、任务调度体系与 PyMongo 驱动衔接起来你不再需要在 DAG 里硬编码数据库地址和密码而是把连接信息集中存放在 Airflow Connection 中在 Python 任务里通过MongoHook获得一个标准的 PyMongo 客户端执行读写用MongoSensor让 DAG 暂停等待某条文档出现后再继续下游任务。读完本文你将掌握该 Provider 的安装方式含 Airflow constraints 固定版本、mongo_default连接的创建与校验、两种最常见的生产工作流代码以及避免调度器/Worker 环境不一致和凭据泄露的实用注意事项。Golden Rule使用前提与核心约定使用apache-airflow-providers-mongo前先建立三条基本认知它不是独立的 MongoDB 客户端必须与固定pinned版本的apache-airflow一同安装它本质上是 Airflow 生态的扩展包运行时依赖 Airflow 核心与pymongo驱动。凭据与连接选项一律放进 Airflow Connection如mongo_default而不是硬编码在 DAG 文件中。这样密钥可以交给 Airflow 的 Secrets Backend 管理DAG 代码也可以在各环境间复用。MongoHook.get_conn()返回的就是一个普通 PyMongo 客户端拿到它之后请直接用标准的 database / collection APIinsert_one、find_one、update_one等不需要学习任何 Airflow 专属的数据库语法。安装与 Airflow 同环境、按 constraints 固定版本Provider 必须安装在运行 DAG 代码的同一套 Python 环境里。官方推荐用 Airflow 的 constraints 文件同时锁定 Airflow 与 Provider 版本避免pip在解析依赖时悄悄改动 Airflow 核心版本python -m venv .venv source .venv/bin/activate python -m pip install --upgrade pip AIRFLOW_VERSIONyour-airflow-version PROVIDER_VERSION5.3.2 PYTHON_VERSION$(python -c import sys; print(f{sys.version_info.major}.{sys.version_info.minor})) CONSTRAINT_URLhttps://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt python -m pip install \ apache-airflow${AIRFLOW_VERSION} \ apache-airflow-providers-mongo${PROVIDER_VERSION} \ --constraint ${CONSTRAINT_URL}如果 Airflow 已经安装好再单独添加 Provider 时同样要把apache-airflow固定住python -m pip install \ apache-airflowyour-airflow-version \ apache-airflow-providers-mongo5.3.2安装后做两个快速自检airflow providers list | grep mongo airflow info第一个命令确认 Provider 已被 Airflow 识别第二个命令检查整体环境信息Python 版本、Airflow 版本、安装路径等。Provider 依赖的 PyMongo 驱动细节可参考仓库内的 PyMongo 驱动指南当前文档对应版本 4.16.0注意它要求 Python 3.9且mongodbsrv://需要dnspython2.6.1。认证与连接设置从环境变量到 Airflow ConnectionProvider 通过 Airflow Connection 读取 MongoDB 凭据。安全的做法是连接值放在环境变量中再基于环境变量创建 Airflow Connection这样密钥不会进入 DAG 代码或版本库。先导出环境变量export MONGO_HOSTmongo.example.com export MONGO_PORT27017 export MONGO_DBanalytics export MONGO_USERairflow export MONGO_PASSWORDsecret再创建 Airflow 连接airflow connections add mongo_default \ --conn-type mongo \ --conn-host $MONGO_HOST \ --conn-port $MONGO_PORT \ --conn-schema $MONGO_DB \ --conn-login $MONGO_USER \ --conn-password $MONGO_PASSWORD各字段语义对照CLI 参数对应连接字段本示例用途--conn-typeconn_type固定为mongo让 Airflow 识别为 MongoDB 连接--conn-hosthostMongoDB 服务器地址--conn-portport端口默认27017--conn-schemaschema数据库名analytics--conn-loginlogin用户名airflow--conn-passwordpassword密码把连接接入 DAG 前先确认它已存在且可读airflow connections get mongo_default如果部署需要 TLS、副本集replica set、SRV 或其它超出 host / port / database / credentials 的连接选项应将这些配置放在 Airflow Connection 本身里通过 connection extra 字段而不是散布在 DAG 代码中。仓库内 MongoDB Atlas 指南 和 PyMongo 驱动指南 对 SRV、Stable API、认证扩展AWS、GSSAPI、OCSP 等有更细的介绍可作为 Connection extra 配置的参考背景。常见工作流一在 Python 任务中读写文档MongoHook当任务需要从 Python 代码里执行常规 MongoDB 操作写入、查询、更新时使用MongoHook。get_conn()返回 PyMongo 的MongoClient因此可以立即使用熟悉的 PyMongo APIfrom __future__ import annotations import pendulum from airflow import DAG from airflow.decorators import task from airflow.providers.mongo.hooks.mongo import MongoHook with DAG( dag_idmongo_hook_example, start_datependulum.datetime(2024, 1, 1, tzUTC), scheduleNone, catchupFalse, tags[mongo], ): task def write_and_read() - None: hook MongoHook(mongo_conn_idmongo_default) client hook.get_conn() try: collection client[analytics][events] collection.insert_one( { event_type: signup, source: airflow, status: queued, } ) document collection.find_one({event_type: signup}) print(document) collection.update_one( {event_type: signup}, {$set: {status: processed}}, ) finally: client.close() write_and_read()这个模式通常是在 Airflow 中操作 MongoDB 的最简路径关键四步用MongoHook(mongo_conn_id...)取得 hook 实例在任务内部只调用一次get_conn()获得客户端对返回的客户端直接使用标准 PyMongo APIcollection.insert_one/find_one/update_one等如果客户端是自己管理的任务退出前记得client.close()释放连接。需要留意的是MongoClient构造本身并不会因凭据错误或服务器不可达而快速失败连接是惰性的因此把client.close()放进finally保证资源回收、并在业务逻辑里对find_one等调用做好异常处理是生产代码的稳妥习惯。更完整的连接验证与超时设置如serverSelectionTimeoutMS可参考 PyMongo 驱动指南。常见工作流二等待匹配文档MongoSensor当 DAG 需要暂停直到某个 collection 中出现符合查询条件的文档时才继续时使用MongoSensor。典型场景是上游系统写入一条“ready”记录下游任务只有看到它才能开始。from __future__ import annotations import pendulum from airflow import DAG from airflow.providers.mongo.sensors.mongo import MongoSensor with DAG( dag_idmongo_sensor_example, start_datependulum.datetime(2024, 1, 1, tzUTC), scheduleNone, catchupFalse, tags[mongo], ): wait_for_ready_document MongoSensor( task_idwait_for_ready_document, mongo_conn_idmongo_default, collectionevents, query{status: ready}, poke_interval30, timeout60 * 20, )MongoSensor的关键参数mongo_conn_id复用与MongoHook相同的 Airflow 连接连接 id 在整个 DAG 中保持一致collection要轮询的目标集合名query匹配条件只要集合中存在任意一条满足该查询的文档Sensor 即成功poke_interval轮询间隔秒上例为 30 秒timeout最长等待时间秒上例为 20 分钟超时后任务失败。性能忠告Sensors 会反复轮询所以查询要尽量小、尽量走索引。一个不设索引的全集合扫描会把简单的“就绪检查”变成数据库上的持续负载。poke_interval越大、query越精确对 MongoDB 的压力就越小。常见配置模式连接、Hook 与 Sensor 的分工对大多数 DAG推荐如下清晰的分工数据库名、主机、凭据→ 全部放在 Airflow Connection如mongo_default里自定义读写/更新逻辑→ 在task函数内使用MongoHook下游任务依赖文档存在→ 使用MongoSensor等待连接 id 跨 DAG 保持稳定→ 例如统一用mongo_default或warehouse_mongo便于运维统一管理、切换环境时只改连接不改代码。Pitfalls六个高频踩坑点Provider 必须安装在所有运行 DAG 代码的位置。Scheduler、Worker 以及本地测试环境只要import airflow.providers.mongo就需要装有该包否则任务会以 ImportError 失败。凭据留在 Airflow Connection 或 Secrets Backend 中。不要在 DAG 代码里直接嵌入 MongoDB 用户名和密码。升级 Provider 时保持apache-airflow固定避免pip静默替换 Airflow 核心版本导致不兼容。使用 Worker 可达的主机名。笔记本上能解析的 MongoDB 主机在容器或远程 Worker 里可能解析不了连接里的 host 要以实际运行任务的网络为准。连接专属选项集中在 Connection 中不要散落在 DAG 代码各处否则切换环境本地 → 生产时难以维护。结合 PyMongo 驱动指南 的版本注意点PyMongo 4.16 要求 Python3.9、mongodbsrv://需要dnspythonProvider 与 MongoDB Server 版本的兼容性彼此独立升级 Airflow 核心后务必复查 Provider 的兼容性与 release notes 再变更生产环境的 pin。版本说明本文覆盖apache-airflow-providers-mongo版本5.3.2Provider 包的版本兼容性是相对 Airflow 而言的与你的 MongoDB Server 版本没有绑定关系。升级 Airflow 核心时应重新核对 Provider 兼容性矩阵与发布说明再修改生产环境的版本 pin。延伸阅读本文聚焦 Mongo Provider 的接入如果你想继续深入仓库内已有这些关联资料可对照阅读Apache Airflow 核心包指南Airflow 本身的安装与使用基础PyMongo 驱动指南MongoHook.get_conn()返回的客户端的完整 API、连接验证与认证扩展MongoDB Atlas 指南托管 MongoDBAtlas场景下的连接方式与最佳实践MongoEngine 指南 与 Motor 指南如需在任务里使用 ODM 或异步驱动可作参考。赞分享【免费下载链接】context-hub项目地址https://gitcode.com/gh_mirrors/co/context-hub点击查看免费下载相关推荐Apache Airflow 集成 AirbyteAirbyte Connection 连接配置完整指南Apache Airflow 集成 AirbyteAirbyte Connection 连接配置完整指南 本篇技术指南聚焦 Apache Airflow 的后端任务调度工作流自动化数据编排批处理数据工程流程编排抖音无水印批量下载douyin-downloader 实操指南抖音无水印批量下载douyin downloader 实操指南 周四下午你在竞品账号的作品列表里点了一下午另存为40 条视频攒出一堆带水印的文件还得网页爬虫CLIApache Airflow Amazon Athena 连接配置实战从 Connection 到 AthenaSQLHook 的完整指南Apache Airflow Amazon Athena 连接配置实战从 Connection 到 AthenaSQLHook 的完整指南 导读 本文围绕 A后端任务调度工作流自动化数据编排批处理数据工程流程编排上一篇终极音频频谱分析指南Spek免费工具让你的音乐可视化下一篇如何用SPT-AKI Profile Editor存档修改器解放你的塔科夫游戏体验创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。