news 2026/10/9 2:14:48

Apache Airflow MySQL Provider 实战指南:连接配置、MySqlHook 与 DAG 数据管道

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow MySQL Provider 实战指南:连接配置、MySqlHook 与 DAG 数据管道

【免费下载链接】context-hub

项目地址:https://gitcode.com/gh_mirrors/co/context-hub
点击查看免费下载

apache-airflow-providers-mysql是 Apache Airflow 的 MySQL Provider 包,它让 DAG 代码通过 Airflow 连接(Connection)与 Hook 机制安全地读写 MySQL。本文以本仓库中 MySQL Provider 指南 为主体,结合仓库内 Airflow 核心包指南 与 common-sql Provider 指南,完整讲解安装、连接配置、MySqlHook常见用法、参数绑定与易错点。读完你可以直接在 Airflow 3 环境中编写可运行的建表、写入、查询 DAG,并避免conn_id与mysql_conn_id混用等高频踩坑。

这个包解决了什么

apache-airflow-providers-mysql是 Airflow 的数据库 Provider 之一,安装在与apache-airflow相同的 Python 环境中,用于让 DAG 通过 Airflow 连接与 Hook 访问 MySQL。

需要注意它的定位:它是为 Airflow 任务与 DAG 代码服务的集成层,不是供普通应用代码使用的独立 MySQL 客户端库。你在 DAG 中应通过MySqlHook或基于该 Provider 的 Operator 访问数据库,而不是在 DAG 里直接 new 一个 pymysql 连接、把凭据硬编码进代码。

在 Context Hub 中,该文档以apache-airflow/providers-mysql(语言变体py,版本6.5.0)为条目 ID 收录,可使用chub get拉取(详见 CLI 参考 与 内容指南)。

安装:与 Airflow 同环境、同版本约束

Provider 必须安装到 Airflow scheduler、worker、webserver 实际运行的那个 Python 环境。官方指南建议锁定 Provider 版本:

python -m pip install "apache-airflow-providers-mysql==6.5.0"

如果你在构建自定义 Airflow 镜像,应在镜像构建阶段加入 Provider,保证每个 Airflow 组件看到的 Provider 集合一致,避免调度器与 Worker 因缺包而导致 import 失败。

结合仓库内 Airflow 核心包指南,更稳妥的完整流程是:先用官方约束文件(constraints file)安装 Airflow 核心,再单独安装 Provider,且安装 Provider 时在同一命令中再次钉住apache-airflow版本,防止 pip 静默升级或降级核心:

AIRFLOW_VERSION=3.1.8 PYTHON_VERSION="$(python -c 'import sys; print(f"{sys.version_info.major}.{sys.version_info.minor}")')" CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt" python -m pip install "apache-airflow==${AIRFLOW_VERSION}" --constraint "${CONSTRAINT_URL}" python -m pip install "apache-airflow==${AIRFLOW_VERSION}" "apache-airflow-providers-mysql==6.5.0"

需要强调:核心包指南明确警告,直接裸执行pip install apache-airflow可能产生不可用的安装;Airflow 3 官方支持的安装器是pip与uv。Provider 的versions字段指 PyPI 上的包版本(本仓库记录为6.5.0),升级 Airflow 核心或 Provider 之前,先检查两者兼容性要求。

配置 MySQL 连接

Provider 从 Airflow 连接中读取凭据与连接设置。默认连接 ID 是mysql_default。

在 Airflow UI 中创建

创建连接时填写:

字段值
Connection Idmysql_default
Connection Typemysql
Host你的 MySQL 主机名
Schema默认数据库名
Login / Password数据库账号与密码
Port通常为3306

用环境变量定义连接

export AIRFLOW_CONN_MYSQL_DEFAULT='mysql://airflow:airflow@mysql:3306/analytics?charset=utf8mb4'

环境变量方式非常适合 CI 与 Agent 场景:仓库 CLI 参考 与核心包指南都强调环境变量是覆盖配置最安全的手段。注意:如果密码中包含@、:、/等 URL 保留字符,必须先做 URL 编码再放入 URI,否则连接解析会出错。

连接 Extras

MySQL 部署常用的连接 extras 如下:

{ "charset": "utf8mb4", "cursor": "dictcursor", "local_infile": false }

extras 用于表达连接级行为,例如字符集、游标类型、SSL 设置、Unix socket 设置、LOAD DATA LOCAL INFILE支持等,而不是在每个 DAG 里硬编码这些细节。其中cursor: "dictcursor"直接决定查询结果的返回形态(见下文“查询参数与结果形状”)。

最小连接检查:用 MySqlHook 做连通性探针

在任务里用MySqlHook快速验证连接是否可用:

from airflow.decorators import task from airflow.providers.mysql.hooks.mysql import MySqlHook @task def ping_mysql() -> int: hook = MySqlHook(mysql_conn_id="mysql_default") row = hook.get_first("SELECT 1") return int(row[0])

关键差异:MySqlHook使用mysql_conn_id而不是conn_id。后者是 Airflow 通用 SQL 操作符(如SQLExecuteQueryOperator)使用的参数名——这是 DAG 作者最容易混淆的一对参数。

常见 DAG 工作流:建表、写入、查询

对大多数 DAG,用MySqlHook完成建表、写行、读结果即可:

from __future__ import annotations from datetime import datetime from airflow.decorators import dag, task from airflow.providers.mysql.hooks.mysql import MySqlHook @dag( dag_id="mysql_provider_example", start_date=datetime(2024, 1, 1), schedule=None, catchup=False, tags=["mysql"], ) def mysql_provider_example(): @task def create_table() -> None: hook = MySqlHook(mysql_conn_id="mysql_default") hook.run( """ CREATE TABLE IF NOT EXISTS events ( id INT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(255) NOT NULL, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) """ ) @task def insert_event() -> None: hook = MySqlHook(mysql_conn_id="mysql_default") hook.run( "INSERT INTO events (name) VALUES (%s)", parameters=("signup",), ) @task def fetch_recent() -> list[tuple[int, str]]: hook = MySqlHook(mysql_conn_id="mysql_default") return hook.get_records( "SELECT id, name FROM events ORDER BY id DESC LIMIT 10" ) create_table() >> insert_event() >> fetch_recent() mysql_provider_example()

这段示例也符合 Airflow 3 的 DAG 编写规范:使用airflow.decorators的@dag/@task(核心包指南指出 Airflow 3 稳定编写接口在airflow.sdk,应避免from airflow.models import DAG等旧路径),显式设置catchup=False。

最常用的 Hook 方法:

方法用途
get_first(sql, parameters=None)取单行结果
get_records(sql, parameters=None)取多行结果
run(sql, parameters=None)执行 DDL 或 DML 语句
get_conn()获取原始 DB-API 连接,手动控制游标

原始游标访问:需要显式控制时再降级

当需要显式游标处理(例如逐行 UPDATE 并手动提交)时,可以降到驱动连接层:

from airflow.providers.mysql.hooks.mysql import MySqlHook hook = MySqlHook(mysql_conn_id="mysql_default") conn = hook.get_conn() cursor = conn.cursor() try: cursor.execute( "UPDATE events SET name=%s WHERE id=%s", ("activated", 1), ) conn.commit() finally: cursor.close() conn.close()

使用get_conn()直接操作时,写操作必须自己 commit。这是与hook.run()的重要区别:run()内部的提交策略由 Hook 管理,而原始连接把事务控制完全交给你。

查询参数与结果形状

MySQL 使用%s占位符绑定参数,不要使用 SQLite 风格的?:

hook.run( "INSERT INTO events (name) VALUES (%s)", parameters=("signup",), )

当希望行以列名为键(而不是位置元组)返回时,在连接 extras 中配置{"cursor": "dictcursor"},get_records等方法的返回项就会按列名索引。

与 common-sql Provider 的关系

仓库内 common-sql Provider 指南 说明:apache-airflow-providers-common-sql提供跨数据库复用的 SQL 操作符(如SQLExecuteQueryOperator)与DbApiHook辅助方法,但数据库连接类型本身来自对应 Provider——MySQL 场景下,即使你使用 common-sql 的通用操作符,也必须安装apache-airflow-providers-mysql,且这些通用操作符使用conn_id参数(如conn_id="warehouse"),与MySqlHook的mysql_conn_id命名不同。同一个 MySQL 连接既可以被MySqlHook(mysql_conn_id=...)使用,也可以被 common-sql 操作符按conn_id引用。

常见陷阱清单

  • MySqlHook期望mysql_conn_id;Airflow 的通用 SQL 操作符使用conn_id——混用会导致连接找不到或参数名报错。
  • 查询参数使用%s占位符,不要使用?。
  • 从get_conn()拿到原始连接后,写操作必须显式 commit。
  • 不要把大型 MySQL 结果集塞进 XCom;保持任务返回值足够小(XCom 会序列化并传输返回值,大结果集会拖垮调度与 Worker 内存)。
  • 只有确实需要本地文件加载且服务器端已配置支持时,才启用local_infile(extras 中"local_infile": true),否则保持关闭以降低安全风险。

版本说明

本指南针对apache-airflow-providers-mysql版本6.5.0(见文档 frontmatter 的versions字段)。保持 Provider 与 DAG 代码运行在同一 Airflow 环境,并在独立升级 Airflow 核心或 Provider 前核对兼容性要求。此外,若将 MySQL 用作 Airflow 自身的元数据库,核心包指南给出的连接串形如mysql+mysqldb://airflow:secret@localhost:3306/airflow(通过AIRFLOW__DATABASE__SQL_ALCHEMY_CONN配置),并强调生产环境应将 SQLite 替换为 PostgreSQL 或 MySQL、且 MariaDB 不被官方支持——这是同一 Provider 在 Airflow 部署层面的另一个典型应用场景。

【免费下载链接】context-hub

项目地址:https://gitcode.com/gh_mirrors/co/context-hub
点击查看免费下载
上一篇:Cursor Free VIP破解工具:3步解决Cursor AI试用限制,永久免费使用Pro功能
下一篇:OmenSuperHub深度解析:惠普游戏本硬件控制与性能调优实战指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/9 2:14:25

HTTP协议零基础拆解:请求头、响应状态码与调试实战

1. 从一次浏览器地址栏输入开始说起如果你正在学 Web 开发,无论你打算写前端、后端、还是做全栈,HTTP 都是那个绕不开的坎。它就像网络世界的普通话,前端和后端沟通、浏览器和服务器沟通、App 和云服务沟通,全都靠它。很多新手被 …

作者头像 李华