
Apache Airflow ArangoDB Provider 集成实战连接配置、AQL 执行与数据感知【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 ArangoDB Providerapache-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 暴露了三类能力Hookairflow.providers.arangodb.hooks.arangodb.ArangoDBHook负责建立并复用 ArangoDB 连接Operatorairflow.providers.arangodb.operators.arangodb包含执行 AQL 的AQLOperator和集合文档操作的ArangoDBCollectionOperatorSensorairflow.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-airflow2.11.0apache-airflow-providers-common-compat1.10.1python-arango7.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.14pyproject.toml中requires-python 3.10与之对应。三、配置 ArangoDB ConnectionArangoDB 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:8530Database/Schema必填目标数据库名如_systemUsername必填用户名如rootPassword必填对应密码注意该连接类型的 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 场景缺失时抛出AirflowExceptiondatabase读取conn.schemausername读取conn.loginpassword读取conn.password允许为空字符串。Hook 默认的连接 ID 为arangodb_defaultdefault_conn_nameconn_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连接复用与基础 APIArangoDBHookhooks/arangodb.py继承自common-compat的BaseHook职责是“连接到 ArangoDB 并取得客户端”是 Operator 与 Sensor 的公共底座。4.1 连接建立机制Hook 通过三个cached_property完成连接对象的惰性初始化与复用cached_property def client(self) - ArangoDBClient: return ArangoDBClient(hostsself.hosts) cached_property def db_conn(self) - StandardDatabase: return self.client.db(nameself.database, usernameself.username, passwordself.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.pydef 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会透传给驱动例如countTrue可让游标携带结果总数AQLSensor 正是依赖这一点。4.3 集合与文档管理方法Hook 还提供了一批幂等性良好的元数据与文档操作方法源码见 hooks/arangodb.py方法行为create_collection(name)集合不存在则创建并返回True已存在则返回Falsedelete_collection(name)集合存在则删除并返回True否则返回Falsecreate_database(name)数据库不存在则创建create_graph(name)图不存在则创建支持 ArangoDB 图模型insert_documents(collection_name, documents)集合不存在时先自动创建再批量插入insert_manyupdate_documents(collection_name, documents)集合不存在时抛AirflowException否则批量更新update_manyreplace_documents(collection_name, documents)批量整体替换replace_many集合不存在时报错delete_documents(collection_name, documents)按_key等条件批量删除delete_many批量操作均捕获python-arango的DocumentInsertError/DocumentUpdateError/DocumentReplaceError/DocumentDeleteError记录错误日志后重新抛出方便在 Airflow 日志中排查。五、AQLOperator在 DAG 中执行 AQLAQLOperatoroperators/arangodb.py用于在 ArangoDB 中执行一条 AQL 查询AQLOperator( task_idaql_operator, queryFOR doc IN students RETURN doc, dagdag, result_processorlambda cursor: print([document[name] for document in cursor]), )参数说明query必填AQL 语句字符串或指向包含 AQL 的.sql文件路径见下文模板机制arangodb_conn_id使用的连接 ID默认arangodb_defaultresult_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_idarangodb_default构造 Hook 并恰好调用一次query()。六、ArangoDBCollectionOperator集合级 CRUDArangoDBCollectionOperatoroperators/arangodb.py面向“集合文档”的常规数据操作构造参数包括arangodb_conn_id连接 ID默认arangodb_defaultcollection_name目标集合名documents_to_insert/documents_to_update/documents_to_replace/documents_to_delete均为 Python dict 列表分别对应插入、更新、替换、删除delete_collection布尔值为True时删除整个集合。典型用法ArangoDBCollectionOperator( task_idinsert_students, collection_namestudents, documents_to_insert[ {_key: lola, first: Lola, last: Martin}, ], )execute()的关键行为源码 operators/arangodb.py若插入、更新、替换、删除、删集合五个操作全部未指定抛出ValueError(At least one operation must be specified.)防止空任务误执行各操作按“插入 → 更新 → 替换 → 删除 → 删除集合”顺序依次执行并在日志中记录操作文档数量。该行为被单元测试覆盖test_arangodb.pytest_insert_documents验证插入路径正确调用insert_documentstest_no_operation_fails验证未指定任何操作时抛出ValueError且不会误调用插入。七、AQLSensor等待 ArangoDB 数据就绪AQLSensorsensors/arangodb.py继承自BaseSensorOperator作用是轮询执行一条 AQL 查询直到查询返回至少一条记录AQLSensor( task_idaql_sensor, queryFOR doc IN students FILTER doc.name judy RETURN doc, timeout60, poke_interval10, dagdag, )query必填AQL 语句或.sql模板文件arangodb_conn_id连接 ID默认arangodb_default继承自 Sensor 基类的timeout超时秒数与poke_interval轮询间隔秒数控制探测节奏。poke()的实现sensors/arangodb.py值得关注它调用hook.query(self.query, countTrue).count()利用countTrue让游标附带总记录数随后返回records ! 0——即记录数非零即视为条件满足。因此该 Sensor 适合“等待某类数据写入后继续下游”例如等待students集合中出现名为judy的文档。单元测试 test_arangodb.py 验证了query().count()会被调用且传感器能正常完成执行。八、完整示例 DAG 与 SQL 模板机制仓库内置的示例 DAG example_arangodb.py 完整展示了上述组件的组合用法dag DAG( example_arangodb_operator, start_datedatetime(2021, 1, 1), tags[example], catchupFalse, ) sensor AQLSensor( task_idaql_sensor, queryFOR doc IN students FILTER doc.name judy RETURN doc, timeout60, poke_interval10, dagdag, ) operator AQLOperator( task_idaql_operator, queryFOR doc IN students RETURN doc, dagdag, result_processorlambda 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_idaql_sensor_template_file, querysearch_judy.sql, timeout60, poke_interval10, dagdag, ) operator2 AQLOperator( task_idaql_operator_template_file, dagdag, result_processorlambda cursor: print([document[name] for document in cursor]), querysearch_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通过countTrue的游标统计记录数非零即满足配合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),仅供参考