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 的 ArangoDB Provider(apache-airflow-providers-arangodb)是 Airflow 官方维护的 ArangoDB 集成包,用于在 DAG 中通过 Hook、Operator 与 Sensor 与 ArangoDB 数据库交互,执行 AQL 查询、管理集合文档并感知数据就绪状态。本文基于当前仓库中的 provider 包说明 及其配套文档与源码,完整讲解该 Provider 的安装依赖、连接配置、核心 Hook API、两大 Operator 与 AQLSensor 的用法,并结合源码与单元测试揭示其底层实现机制。
一、Provider 包概览与适用场景
apache-airflow-providers-arangodb是 Apache Airflow 的一个 Provider 发行包(release 版本为2.9.6),其所有类都位于airflow.providers.arangodbPython 包中。该包在provider.yaml中声明为state: ready、lifecycle: production,说明它处于生产可用状态,源码见 provider.yaml。
从provider.yaml中的模块注册信息可以看到,该 Provider 向 Airflow 暴露了三类能力:
- Hook:
airflow.providers.arangodb.hooks.arangodb.ArangoDBHook,负责建立并复用 ArangoDB 连接; - Operator:
airflow.providers.arangodb.operators.arangodb,包含执行 AQL 的AQLOperator和集合文档操作的ArangoDBCollectionOperator; - Sensor:
airflow.providers.arangodb.sensors.arangodb.AQLSensor,轮询等待符合条件的文档出现。
典型应用场景包括:ETL/ELT 流水线中向 ArangoDB 写入清洗后的文档、周期性地用 AQL 聚合图数据、以及在依赖下游任务前先等待 ArangoDB 中某批数据写入完成(数据就绪感知)。
二、安装与依赖要求
2.1 pip 安装
在已有 Airflow 环境基础上,直接通过 pip 安装:
pip install apache-airflow-providers-arangodb2.2 版本要求
根据 README.rst 与 pyproject.toml,该 Provider 的依赖要求如下:
| PIP 包 | 版本要求 |
|---|---|
apache-airflow | >=2.11.0 |
apache-airflow-providers-common-compat | >=1.10.1 |
python-arango | >=7.3.2 |
其中python-arango是 ArangoDB 官方 Python 驱动,Hook 的底层连接与 AQL 执行全部基于它完成;apache-airflow-providers-common-compat提供跨 Airflow 版本的兼容抽象(BaseHook、BaseOperator、BaseSensorOperator等)。该包要求 Python 版本为 3.10、3.11、3.12、3.13、3.14,pyproject.toml中requires-python = ">=3.10"与之对应。
三、配置 ArangoDB Connection
ArangoDB Provider 通过 Airflow 标准的 Connection 机制管理凭据,连接类型为arangodb,参考文档见 connections/arangodb.rst。
3.1 字段说明
| 字段 | 是否必填 | 说明与示例 |
|---|---|---|
| Host | 必填 | ArangoDB 主机 URL,支持逗号分隔的多个 URL(集群场景下的多个 coordinator),如http://127.0.0.1:8529或http://127.0.0.1:8529,http://127.0.0.1:8530 |
| Database/Schema | 必填 | 目标数据库名,如_system |
| Username | 必填 | 用户名,如root |
| Password | 必填 | 对应密码 |
注意:该连接类型的 UI 表单会隐藏Port与Extra字段(端口已包含在 Host URL 中),并重新标记字段含义,详见 provider.yaml 中connection-types的ui-field-behaviour声明。
3.2 源码中的字段解析逻辑
在 hooks/arangodb.py 中,ArangoDBHook通过四个属性把 Connection 映射到 ArangoDB 客户端参数:
hosts:读取conn.host并用","切分为 URL 列表(对应集群多 coordinator 场景),缺失时抛出AirflowException;database:读取conn.schema;username:读取conn.login;password:读取conn.password(允许为空字符串)。
Hook 默认的连接 ID 为arangodb_default(default_conn_name),conn_name_attr = "arangodb_conn_id"则允许 DAG 中通过arangodb_conn_id参数指定任意自定义连接。
3.3 用 CLI 创建连接
不依赖 UI 时,可以用 Airflow CLI 等价创建(host 中带端口,schema 填数据库名):
airflow connections add arangodb_default \ --conn-type arangodb \ --conn-host 'http://127.0.0.1:8529' \ --conn-login root \ --conn-password 'your_password' \ --conn-schema _system四、ArangoDBHook:连接复用与基础 API
ArangoDBHook(hooks/arangodb.py)继承自common-compat的BaseHook,职责是“连接到 ArangoDB 并取得客户端”,是 Operator 与 Sensor 的公共底座。
4.1 连接建立机制
Hook 通过三个cached_property完成连接对象的惰性初始化与复用:
@cached_property def client(self) -> ArangoDBClient: return ArangoDBClient(hosts=self.hosts) @cached_property def db_conn(self) -> StandardDatabase: return self.client.db(name=self.database, username=self.username, password=self.password) @cached_property def _conn(self) -> Connection: return self.get_connection(self.arangodb_conn_id)调用链为:get_connection()从 Airflow 元数据库加载 Connection → 用hosts构造python-arango的ArangoClient→ 再以database/username/password取得目标数据库的StandardDatabase包装对象。cached_property保证同一 Hook 实例内客户端与数据库句柄只创建一次,避免每个任务重复握手。get_conn()方法返回缓存的client,供自定义扩展使用。
4.2 AQL 查询
query()方法是整个 Provider 的核心入口(hooks/arangodb.py):
def query(self, query, **kwargs) -> Cursor: result = self.db_conn.aql.execute(query, **kwargs) if not isinstance(result, Cursor): raise AirflowException("Failed to execute AQLQuery, expected result to be of type Cursor") return result它通过db_conn.aql.execute()执行 AQL 语句并返回arango.cursor.Cursor游标;若返回类型不符或执行出错(AQLQueryExecuteError),统一包装为AirflowException抛出。**kwargs会透传给驱动,例如count=True可让游标携带结果总数(AQLSensor 正是依赖这一点)。
4.3 集合与文档管理方法
Hook 还提供了一批幂等性良好的元数据与文档操作方法(源码见 hooks/arangodb.py):
| 方法 | 行为 |
|---|---|
create_collection(name) | 集合不存在则创建并返回True,已存在则返回False |
delete_collection(name) | 集合存在则删除并返回True,否则返回False |
create_database(name) | 数据库不存在则创建 |
create_graph(name) | 图不存在则创建(支持 ArangoDB 图模型) |
insert_documents(collection_name, documents) | 集合不存在时先自动创建,再批量插入(insert_many) |
update_documents(collection_name, documents) | 集合不存在时抛AirflowException,否则批量更新(update_many) |
replace_documents(collection_name, documents) | 批量整体替换(replace_many),集合不存在时报错 |
delete_documents(collection_name, documents) | 按_key等条件批量删除(delete_many) |
批量操作均捕获python-arango的DocumentInsertError/DocumentUpdateError/DocumentReplaceError/DocumentDeleteError,记录错误日志后重新抛出,方便在 Airflow 日志中排查。
五、AQLOperator:在 DAG 中执行 AQL
AQLOperator(operators/arangodb.py)用于在 ArangoDB 中执行一条 AQL 查询:
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]), )参数说明:
query:必填,AQL 语句字符串,或指向包含 AQL 的.sql文件路径(见下文模板机制);arangodb_conn_id:使用的连接 ID,默认"arangodb_default";result_processor:可选的Callable,接收查询返回的Cursor并做进一步处理,例如逐文档打印、写入下游存储或做校验;template_fields = ("query",):query参与 Airflow 模板渲染,可在其中使用{{ ds }}、{{ params.xxx }}等 Jinja 变量。
execute()内部流程(operators/arangodb.py):实例化ArangoDBHook(arangodb_conn_id=...)→ 调用hook.query(query)得到游标 → 若提供了result_processor则把游标传入执行。单元测试 test_arangodb.py 验证了AQLOperator会以arangodb_conn_id="arangodb_default"构造 Hook 并恰好调用一次query()。
六、ArangoDBCollectionOperator:集合级 CRUD
ArangoDBCollectionOperator(operators/arangodb.py)面向“集合+文档”的常规数据操作,构造参数包括:
arangodb_conn_id:连接 ID,默认"arangodb_default";collection_name:目标集合名;documents_to_insert/documents_to_update/documents_to_replace/documents_to_delete:均为 Python dict 列表,分别对应插入、更新、替换、删除;delete_collection:布尔值,为True时删除整个集合。
典型用法:
ArangoDBCollectionOperator( task_id="insert_students", collection_name="students", documents_to_insert=[ {"_key": "lola", "first": "Lola", "last": "Martin"}, ], )execute()的关键行为(源码 operators/arangodb.py):
- 若插入、更新、替换、删除、删集合五个操作全部未指定,抛出
ValueError("At least one operation must be specified."),防止空任务误执行; - 各操作按“插入 → 更新 → 替换 → 删除 → 删除集合”顺序依次执行,并在日志中记录操作文档数量。
该行为被单元测试覆盖(test_arangodb.py):test_insert_documents验证插入路径正确调用insert_documents;test_no_operation_fails验证未指定任何操作时抛出ValueError且不会误调用插入。
七、AQLSensor:等待 ArangoDB 数据就绪
AQLSensor(sensors/arangodb.py)继承自BaseSensorOperator,作用是轮询执行一条 AQL 查询,直到查询返回至少一条记录:
AQLSensor( task_id="aql_sensor", query="FOR doc IN students FILTER doc.name == 'judy' RETURN doc", timeout=60, poke_interval=10, dag=dag, )query:必填,AQL 语句或.sql模板文件;arangodb_conn_id:连接 ID,默认"arangodb_default";- 继承自 Sensor 基类的
timeout(超时秒数)与poke_interval(轮询间隔秒数)控制探测节奏。
poke()的实现(sensors/arangodb.py)值得关注:它调用hook.query(self.query, count=True).count(),利用count=True让游标附带总记录数,随后返回records != 0——即记录数非零即视为条件满足。因此该 Sensor 适合“等待某类数据写入后继续下游”,例如等待students集合中出现名为judy的文档。单元测试 test_arangodb.py 验证了query().count()会被调用且传感器能正常完成执行。
八、完整示例 DAG 与 SQL 模板机制
仓库内置的示例 DAG example_arangodb.py 完整展示了上述组件的组合用法:
dag = DAG( "example_arangodb_operator", start_date=datetime(2021, 1, 1), tags=["example"], catchup=False, ) sensor = AQLSensor( task_id="aql_sensor", query="FOR doc IN students FILTER doc.name == 'judy' RETURN doc", timeout=60, poke_interval=10, dag=dag, ) 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]), )8.1 使用.sql模板文件
AQLOperator与AQLSensor都声明了template_ext = (".sql",)与template_fields = ("query",)(见 operators/arangodb.py 与 sensors/arangodb.py),这意味着query参数可以是.sql文件名,Airflow 会自动从dags/ 目录加载该文件内容作为查询语句,并支持 Jinja 渲染:
sensor2 = AQLSensor( task_id="aql_sensor_template_file", query="search_judy.sql", timeout=60, poke_interval=10, dag=dag, ) 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", )若查询文件不在默认的dags/目录,需要在创建 DAG 时通过template_searchpath显式指定搜索路径,示例 DAG 与操作指南 operators/index.rst 中对此有明确说明。
九、源码级运行链路小结
将以上模块串起来,一条完整的数据写入 + 感知 + 消费链路如下:
- 连接:任何组件的
execute()/poke()都会以arangodb_conn_id实例化ArangoDBHook; - 建连:Hook 的
cached_property依次解析hosts → client → db_conn,底层委托python-arango的ArangoClient.db(name, username, password); - 执行:
AQLOperator/ArangoDBCollectionOperator/AQLSensor分别调用query()或集合文档方法,所有驱动异常统一收敛为AirflowException,保证失败任务能被 Airflow 正确重试与告警; - 感知:
AQLSensor通过count=True的游标统计记录数,非零即满足,配合timeout/poke_interval控制轮询行为。
该链路同时被 hooks 测试、operators 测试 与 sensors 测试 三套单元测试覆盖,可作为二次开发或排查问题的参考起点。如需扩展自定义逻辑,AQLOperator的result_processor回调与ArangoDBHook的get_conn()都是官方预留的扩展点。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考