Apache Airflow ArangoDB Provider 实战指南:连接配置、AQL 算子与文档传感器
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 的arangodbProvider 为工作流调度平台与多模型数据库 ArangoDB 之间提供了官方集成能力。本文以 providers/arangodb/docs/index.rst 为骨架,结合仓库内的 Hook、Operator、Sensor 源码与单元测试,系统讲解该 Provider 的安装要求、连接配置、AQL 查询执行、模板化查询文件与集合文档操作,帮助你直接在 DAG 中编排 ArangoDB 的数据读写与条件等待任务。
读完本文,你将掌握:如何用 Airflow Connection 承载 ArangoDB 集群连接参数;如何使用AQLOperator执行并后处理 AQL 查询;如何使用AQLSensor等待目标文档出现;以及如何通过ArangoDBCollectionOperator批量完成集合文档的增改删。
一、Provider 概览与安装要求
apache-airflow-providers-arangodb是 Airflow 官方维护的 Provider 发行包,其状态在 provider.yaml 中标记为ready、生命周期为production,说明该集成已经过生产级成熟度验证。包内所有类均位于 Python 包airflow.providers.arangodb下,由三个核心模块组成:
- hooks:
airflow.providers.arangodb.hooks.arangodb.ArangoDBHook,负责建立与 ArangoDB 的连接并封装底层操作; - operators:
airflow.providers.arangodb.operators.arangodb,包含执行 AQL 的AQLOperator与执行集合操作的ArangoDBCollectionOperator; - sensors:
airflow.providers.arangodb.sensors.arangodb.AQLSensor,用于轮询等待满足条件的文档出现。
1.1 安装方式
在已有 Airflow 环境之上,通过 pip 安装即可:
pip install apache-airflow-providers-arangodb当前仓库中该 Provider 的版本为2.9.6(见 index.rst 与 provider.yaml 的versions列表)。
1.2 版本依赖要求
根据 index.rst 中 "Requirements" 一节的声明,本 Provider 发行版对依赖有如下最低版本要求:
| PIP 包 | 要求版本 |
|---|---|
apache-airflow | >=2.11.0 |
apache-airflow-providers-common-compat | >=1.10.1 |
python-arango | >=7.3.2 |
- apache-airflow >= 2.11.0:即 Airflow 核心的最低支持版本;
- python-arango >= 7.3.2:底层 ArangoDB 官方 Python 驱动,Hook 直接基于
arango.ArangoClient构建(见 hooks/arangodb.py); - apache-airflow-providers-common-compat:提供跨版本兼容的
BaseHook、BaseOperator、BaseSensorOperator与AirflowException等基类,使 Provider 可以同时适配新旧 Airflow 版本。
注意:Provider 中
provider.yaml、docs/下的版本号由发布流程自动维护,手工更新仅限特殊情况;安装时应以 PyPI 上最新发布的 wheel/sdist 为准。
1.3 官方发行包下载与校验
如需在离线环境部署,可以从 Apache 官方下载站获取发行包,并核对签名与校验和:
- sdist 包:
apache_airflow_providers_arangodb-2.9.6.tar.gz(配套.asc签名与.sha512校验文件); - wheel 包:
apache_airflow_providers_arangodb-2.9.6-py3-none-any.whl(同样配套.asc与.sha512)。
从源码构建安装的方式可参考 installing-providers-from-sources.rst。
二、ArangoDB Connection 连接配置详解
ArangoDB Provider 复用了 Airflow 的 Connection 机制来管理连接凭据。官方指南 connections/arangodb.rst 明确指出:ArangoDB 连接提供访问 ArangoDB 所需的凭据,其conn_type为arangodb。
2.1 连接字段与映射关系
从 hooks/arangodb.py 中ArangoDBHook的get_ui_field_behaviour方法可以看到,UI 表单会隐藏port与extra两个字段,并把标准 Connection 字段重新标注为:
| Connection 字段 | UI 中的含义 | 是否必填 | 示例 |
|---|---|---|---|
Host | ArangoDB Host URL,或集群中多个 coordinator 的逗号分隔 URL 列表 | 必填 | http://127.0.0.1:8529或http://127.0.0.1:8529,http://127.0.0.1:8530 |
Schema(Database/Schema) | ArangoDB 数据库名 | 必填 | _system |
Login(Username) | ArangoDB 用户名 | 必填 | root |
Password | ArangoDB 密码 | 必填 | password |
Port | 被 UI 隐藏 | — | — |
Extra | 被 UI 隐藏 | — | — |
值得注意的细节:由于该连接类型把多个 coordinator 地址以逗号分隔写入Host字段,因此Port与Extra字段被刻意隐藏,避免用户填入与hosts解析逻辑冲突的信息。
2.2 Hook 是如何解析连接参数的
ArangoDBHook中定义了hosts、database、username、password四个属性来解析 Connection:
hosts:读取_conn.host后以,分割成列表返回,这就是ArangoClient(hosts=...)接收的多 coordinator 地址来源;database:读取_conn.schema(注意用的是 Schema 字段承载数据库名);username:读取_conn.login;password:读取_conn.password,允许为空字符串。
如果缺少必填的 Host 或 Database 或 Username,Hook 会抛出AirflowException并给出明确的提示信息,例如:
raise AirflowException(f"No ArangoDB Host(s) provided in connection: {self.arangodb_conn_id!r}.")在测试 sensors/test_arangodb.py 中可以看到一个完整的最小连接示例:
Connection( conn_id="arangodb_default", conn_type="arangodb", host="http://127.0.0.1:8529", login="root", password="password", schema="_system", )2.3 连接对象与客户端缓存
ArangoDBHook继承自BaseHook,默认连接 ID 为arangodb_default(default_conn_name),连接类型标识为arangodb。其内部通过cached_property缓存了三类对象:
client:ArangoClient(hosts=self.hosts),即底层 python-arango 客户端;db_conn:self.client.db(name=self.database, username=self.username, password=self.password),即StandardDatabase数据库 API 封装;_conn:get_connection(self.arangodb_conn_id)返回的 Airflow Connection 对象。
这意味着同一任务内重复调用 Hook 不会重复创建连接;而get_conn()方法对外返回的就是缓存的client。
三、AQLOperator:在 DAG 中执行 AQL 查询
操作指南 operators/index.rst 指出:使用AQLOperator可以在 ArangoDB 中执行 AQL 查询,其实现位于 operators/arangodb.py。
3.1 核心参数
| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
query | str | 要执行的 AQL 语句;也可以传一个.sql模板文件路径 | 必填 |
arangodb_conn_id | str | 引用的 ArangoDB Connection ID | arangodb_default |
result_processor | Callable | 对查询结果(Cursor)进一步处理的回调函数 | None |
从源码可以看到,query被声明为template_fields,template_ext为(".sql",),template_fields_renderers将query渲染器指定为"sql"——这意味着查询字符串支持 Jinja 模板渲染,并且可以从.sql文件中加载。
execute方法的核心调用链为:
hook = ArangoDBHook(arangodb_conn_id=self.arangodb_conn_id) result = hook.query(self.query) if self.result_processor: self.result_processor(result)即:实例化 Hook → 调用hook.query()执行 AQL → 若提供了result_processor则把Cursor结果交由其处理。
3.2 实际示例:列出 students 集合的全部文档
官方示例 DAG example_arangodb.py 中的[START howto_aql_operator_arangodb]片段如下:
operator = AQLOperator( task_id="aql_operator", query="FOR doc IN students RETURN doc", dag=dag, result_processor=lambda cursor: print([document["name"] for document in cursor]), )该示例演示了两种能力:
- 用一条原生 AQL 查询遍历
students集合并返回全部文档; - 通过
result_processor回调消费 ArangoDB 返回的Cursor,在任务内直接打印所有文档的name字段——这为后续将结果写入下游任务或外部系统提供了扩展点。
3.3 用 .sql 模板文件加载查询
官方文档特别说明:也可以提供.sql文件来加载查询。文件路径默认相对于dags/目录;如果需要使用其他路径,必须在创建 DAG 对象时提供template_searchpath。
示例 DAG 中的[START howto_aql_operator_template_file_arangodb]片段:
operator2 = AQLOperator( task_id="aql_operator_template_file", dag=dag, result_processor=lambda cursor: print([document["name"] for document in cursor]), query="search_all.sql", )对应 DAG 定义(同一文件顶部):
dag = DAG( "example_arangodb_operator", start_date=datetime(2021, 1, 1), tags=["example"], catchup=False, )使用模板文件时的注意点:
- 查询文件内容可以包含 Jinja 变量,例如
FOR doc IN students FILTER doc.name == '{{ target_name }}' RETURN doc,由 Airflow 渲染后执行; - 若
search_all.sql不在dags/目录下,需要在 DAG 中显式设置template_searchpath=["/path/to/sql_files"]; query字段在 Web UI 的日志/详情中会以 SQL 渲染器展示,便于审查。
3.4 单元测试对执行链路的验证
operators/test_arangodb.py 中的TestAQLOperator::test_arangodb_operator_test通过 mock 验证了执行链路:
op = AQLOperator(task_id="basic_aql_task", query=arangodb_query) op.execute(mock.MagicMock()) mock_hook.assert_called_once_with(arangodb_conn_id="arangodb_default") mock_hook.return_value.query.assert_called_once_with(arangodb_query)这从测试层面确认了:AQLOperator使用默认连接arangodb_default实例化 Hook,并把query原样传递给hook.query()。
四、AQLSensor:等待目标文档出现
当工作流需要"等待某条数据就绪"才能继续时,可使用AQLSensor。官方文档描述为:使用 AQL 查询在 ArangoDB 中等待一个文档或集合出现。
4.1 工作原理
AQLSensor 继承自BaseSensorOperator,核心逻辑在poke方法中:
hook = ArangoDBHook(self.arangodb_conn_id) records = hook.query(self.query, count=True).count() self.log.info("Total records found: %d", records) return records != 0它执行 AQL 查询并请求count=True,然后统计返回的记录数:记录数不为 0 即返回True(目标出现),否则继续轮询,直到达到timeout上限。因此 Sensor 的等待判定本质上依赖查询结果是否为空,例如FILTER doc.name == 'judy'找不到匹配文档时返回空集。
4.2 参数说明
AQLSensor构造函数参数:
| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
query | str | 用于探测的 AQL 查询,或.sql文件路径 | 必填 |
arangodb_conn_id | str | 使用的 ArangoDB Connection ID | arangodb_default |
timeout | int | 超时时间(秒),继承自BaseSensorOperator | 由 Sensor 基类控制 |
poke_interval | int | 两次探测之间的间隔(秒) | 由 Sensor 基类控制 |
与AQLOperator相同,query也是模板字段,支持.sql模板文件与 Jinja 渲染。
4.3 示例:等待 students 集合中出现名为 judy 的文档
sensor = AQLSensor( task_id="aql_sensor", query="FOR doc IN students FILTER doc.name == 'judy' RETURN doc", timeout=60, poke_interval=10, dag=dag, )该示例会在 60 秒超时内每 10 秒执行一次 AQL 查询,直到students集合中存在name == 'judy'的文档为止。
4.4 示例:使用模板文件版
sensor2 = AQLSensor( task_id="aql_sensor_template_file", query="search_judy.sql", timeout=60, poke_interval=10, dag=dag, )与 Operator 同理,search_judy.sql默认从dags/目录解析,如需自定义路径要在 DAG 上设置template_searchpath。
4.5 测试证据
传感器测试 sensors/test_arangodb.py 构造了query.return_value.count.return_value = 1的 mock Hook,验证当查询命中 1 条记录时 Sensor 判定为成功;测试连接同样使用arangodb_default。这印证了"非空结果即成功"的判定语义。
五、ArangoDBCollectionOperator:集合文档的批量增改删
除 AQL 查询外,同文件 operators/arangodb.py 还提供ArangoDBCollectionOperator,用于对集合执行结构化文档操作,无需手写 AQL。
5.1 参数一览
| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
arangodb_conn_id | str | ArangoDB Connection ID | arangodb_default |
collection_name | str | 要操作的集合名称 | 必填 |
documents_to_insert | list[dict] | 要插入的文档字典列表 | [] |
documents_to_update | list[dict] | 要更新的文档字典列表 | [] |
documents_to_replace | list[dict] | 要替换的文档字典列表 | [] |
documents_to_delete | list[dict] | 要删除的文档字典列表 | [] |
delete_collection | bool | 若为True,删除整个集合 | False |
约束:上述五种操作必须至少指定一种,否则execute会抛出ValueError("At least one operation must be specified.")——这一点同样被单元测试test_no_operation_fails覆盖验证。
5.2 底层 Hook 方法语义
ArangoDBCollectionOperator.execute按顺序执行对应操作,最终落到ArangoDBHook的文档操作方法(见 hooks/arangodb.py):
insert_documents(collection_name, documents):若集合不存在则先自动创建(create_collection),再调用collection.insert_many(documents, silent=True)批量插入;失败时记录DocumentInsertError日志并向上抛出;update_documents:要求集合必须存在(否则抛AirflowException),执行update_many(..., silent=True);replace_documents:同样要求集合存在,执行replace_many(..., silent=True);delete_documents:要求集合存在,执行delete_many(..., silent=True);delete_collection:若集合存在则删除,返回True;不存在则记录日志返回False。
silent=True意味着这些批量方法不返回每个操作的详细文档,以减少响应开销;所有批量操作失败时都会先记录self.log.error(...)再重新抛出异常,保证任务状态可被 Airflow 正确标记为失败。
5.3 用法示例
from airflow.providers.arangodb.operators.arangodb import ArangoDBCollectionOperator insert_students = ArangoDBCollectionOperator( task_id="insert_students", collection_name="students", documents_to_insert=[ {"_key": "lola", "first": "Lola", "last": "Martin"}, {"_key": "judy", "first": "Judy", "last": "Smith"}, ], )单元测试 test_insert_documents 验证了该调用会转化为hook.insert_documents("students", documents_to_insert)。删除集合时则设置delete_collection=True,其余文档操作参数留空。
六、Hook 的其他实用能力
ArangoDBHook除了支撑上述 Operator 与 Sensor,还对外暴露了若干便于在自定义算子中复用的方法:
query(query, **kwargs):执行 AQL 并返回arango.cursor.Cursor;若底层抛出AQLQueryExecuteError会被包装为AirflowException,且校验返回值必须是Cursor类型;create_collection(name)/delete_collection(name):创建(不存在时)或删除集合;create_database(name):创建不存在的数据库;create_graph(name):创建不存在的图(Graph);insert_documents/update_documents/replace_documents/delete_documents:文档级批量操作。
由于 Hook 继承自BaseHook,还自带get_connection、log等基础设施,并在 provider.yaml 中声明为arangodb连接类型的hook-class-name。
七、在 DAG 中组合使用:从等待到查询的完整流程
结合 Operator 与 Sensor,可以构建一个典型的"数据就绪 → 读取处理"工作流:先用AQLSensor等待students集合中出现judy文档(超时 60 秒、每 10 秒探测一次),再通过AQLOperator读取并打印全部学生姓名,最后用ArangoDBCollectionOperator完成后续数据维护。三者的执行顺序通过sensor >> operator >> collection_operator等依赖关系编排,查询语句既可以内联书写,也可以外置为.sql模板文件并通过template_searchpath指定位置,兼顾可读性与复用性。
八、常见问题与注意事项
- Host 必须包含协议与端口:
ArangoClient需要完整 URL(如http://127.0.0.1:8529),缺少协议前缀会导致连接失败;集群环境使用逗号分隔多个 coordinator URL,例如http://127.0.0.1:8529,http://127.0.0.1:8530。 - 数据库名填在 Schema 字段:ArangoDB Provider 将 Connection 的
Schema字段解释为数据库名,默认示例为_system;不要在Extra或Port中填写这些信息(UI 中这两个字段已被隐藏)。 - 模板文件路径基准是
dags/:.sql查询文件默认从 DAG 目录解析;放在其他目录时必须设置template_searchpath,否则运行时会找不到文件。 - Sensor 的判定依据是结果非空:若查询本身写错(例如字段名拼写错误),可能永远返回空集导致超时失败,请先用
AQLOperator或 ArangoDB 自带的 Web 界面验证查询。 - Collection 操作约束:
ArangoDBCollectionOperator至少需要指定一种操作;update/replace/delete操作要求集合已存在,只有insert会自动创建集合。 - 版本兼容:使用本 Provider 需要 Airflow >= 2.11.0,且依赖
python-arango >= 7.3.2与apache-airflow-providers-common-compat >= 1.10.1,升级时注意保持版本矩阵一致。
九、参考资料与延伸阅读
- Provider 包索引:providers/arangodb/docs/index.rst
- 连接配置指南:providers/arangodb/docs/connections/arangodb.rst
- 算子与传感器指南:providers/arangodb/docs/operators/index.rst
- Hook 实现:providers/arangodb/src/airflow/providers/arangodb/hooks/arangodb.py
- Operator 实现:providers/arangodb/src/airflow/providers/arangodb/operators/arangodb.py
- Sensor 实现:providers/arangodb/src/airflow/providers/arangodb/sensors/arangodb.py
- 示例 DAG:providers/arangodb/src/airflow/providers/arangodb/example_dags/example_arangodb.py
- Provider 元数据:providers/arangodb/provider.yaml
- 单元测试:providers/arangodb/tests/unit/arangodb/operators/test_arangodb.py、providers/arangodb/tests/unit/arangodb/sensors/test_arangodb.py
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考