【免费下载链接】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 Id | mysql_default |
| Connection Type | mysql |
| 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
相关推荐
SQLite TCL 扩展(tclsqlite)安装与测试全指南:tclextension 目标与 TEA(ish) 构建实战
SQLite TCL 扩展(tclsqlite)安装与测试全指南:tclextension 目标与 TEA ish 构建实战 导读 SQLite 的 TCL 扩
MediaCrawler 上手指南:4 条命令跑通多平台自媒体数据采集
MediaCrawler 上手指南:4 条命令跑通多平台自媒体数据采集 MediaCrawler 是一款开源的多平台自媒体数据采集工具,覆盖小红书、抖音、快手、
网页爬虫数据工程Apache Airflow JDBC Provider(apache-airflow-providers-jdbc 5.4.0)实战指南:用 JDBC 驱动连接任意数据库并编排 DAG 任务
Apache Airflow JDBC Provider(apache airflow providers jdbc 5.4.0)实战指南:用 JDBC 驱动连
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考