Apache Airflow ArangoDB Provider 集成实战:连接配置、AQL 执行与数据感知
2026/9/14 19:16:59 网站建设 项目流程

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: readylifecycle: production,说明它处于生产可用状态,源码见 provider.yaml。

provider.yaml中的模块注册信息可以看到,该 Provider 向 Airflow 暴露了三类能力:

  • Hookairflow.providers.arangodb.hooks.arangodb.ArangoDBHook,负责建立并复用 ArangoDB 连接;
  • Operatorairflow.providers.arangodb.operators.arangodb,包含执行 AQL 的AQLOperator和集合文档操作的ArangoDBCollectionOperator
  • Sensorairflow.providers.arangodb.sensors.arangodb.AQLSensor,轮询等待符合条件的文档出现。

典型应用场景包括:ETL/ELT 流水线中向 ArangoDB 写入清洗后的文档、周期性地用 AQL 聚合图数据、以及在依赖下游任务前先等待 ArangoDB 中某批数据写入完成(数据就绪感知)。

二、安装与依赖要求

2.1 pip 安装

在已有 Airflow 环境基础上,直接通过 pip 安装:

pip install apache-airflow-providers-arangodb

2.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 版本的兼容抽象(BaseHookBaseOperatorBaseSensorOperator等)。该包要求 Python 版本为 3.10、3.11、3.12、3.13、3.14,pyproject.tomlrequires-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:8529http://127.0.0.1:8529,http://127.0.0.1:8530
Database/Schema必填目标数据库名,如_system
Username必填用户名,如root
Password必填对应密码

注意:该连接类型的 UI 表单会隐藏PortExtra字段(端口已包含在 Host URL 中),并重新标记字段含义,详见 provider.yaml 中connection-typesui-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_defaultdefault_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-compatBaseHook,职责是“连接到 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-arangoArangoClient→ 再以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-arangoDocumentInsertError/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_documentstest_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模板文件

AQLOperatorAQLSensor都声明了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 中对此有明确说明。

九、源码级运行链路小结

将以上模块串起来,一条完整的数据写入 + 感知 + 消费链路如下:

  1. 连接:任何组件的execute()/poke()都会以arangodb_conn_id实例化ArangoDBHook
  2. 建连:Hook 的cached_property依次解析hosts → client → db_conn,底层委托python-arangoArangoClient.db(name, username, password)
  3. 执行AQLOperator/ArangoDBCollectionOperator/AQLSensor分别调用query()或集合文档方法,所有驱动异常统一收敛为AirflowException,保证失败任务能被 Airflow 正确重试与告警;
  4. 感知AQLSensor通过count=True的游标统计记录数,非零即满足,配合timeout/poke_interval控制轮询行为。

该链路同时被 hooks 测试、operators 测试 与 sensors 测试 三套单元测试覆盖,可作为二次开发或排查问题的参考起点。如需扩展自定义逻辑,AQLOperatorresult_processor回调与ArangoDBHookget_conn()都是官方预留的扩展点。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询