Apache Airflow 与 Amazon Neptune:使用 Start / Stop 运算符管理图数据库集群生命周期
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Amazon Neptune 是无服务器图数据库服务,专为高性能、高可扩展性与高可用性设计,内置安全能力、持续备份以及与其它 AWS 服务的集成。Apache Airflow 的 amazon provider 提供了两个专门运算符NeptuneStartDbClusterOperator与NeptuneStopDbClusterOperator,让你能够以可编程方式启动、停止 Neptune 数据库集群,并自动等待集群达到目标状态。本文基于官方文档 neptune.rst,结合源码、Hook、Trigger、Waiter 配置与测试代码,全面讲解这两个运算符的使用方法、参数语义、可延迟(deferrable)执行模式以及底层实现原理,帮助你直接在 DAG 中安全、高效地管理 Neptune 集群生命周期。
前置准备
1. 安装 amazon provider
pip install 'apache-airflow[amazon]'详细安装说明请参考 Airflow 安装指南。
2. 准备 AWS 资源与凭据
使用运算符前,需先在 AWS Console 或 AWS CLI 创建必要的 Neptune 集群资源,并配置 Airflow 的 AWS 连接。连接配置详见 AWS 连接指南。
注意:运算符只对已存在的 Neptune 数据库集群执行启动/停止操作,不会创建集群。集群的创建、删除等管理动作,需要你在 AWS 侧自行完成(例如通过 system test 中的
create_db_cluster调用,见后文)。
通用参数
Neptune 运算符继承自AwsBaseOperator,因此支持一系列通用 AWS 参数,这些参数在 generic_parameters.rst 中有完整说明:
| 参数 | 说明 | 默认值 |
|---|---|---|
aws_conn_id | AWS 连接 ID。若设为None,则使用默认的 boto3 行为(不进行连接查询);否则使用连接中存储的凭据 | aws_default |
region_name | AWS 区域名。若为None,使用 AWS 连接 Extra 参数中的 region_name;否则覆盖连接值 | None |
verify | 是否校验 SSL 证书。False表示不校验;也可指定 CA 证书 bundle 文件路径。若为None,使用连接 Extra 参数中的 verify | None |
botocore_config | 用于构造botocore.config.Config的字典,可配置重试策略、超时等 | None |
botocore_config示例(用于配置重试与超时):
{ "signature_version": "unsigned", "s3": { "us_east_1_regional_endpoint": True, }, "retries": { "mode": "standard", "max_attempts": 10, }, "connect_timeout": 300, "read_timeout": 300, "tcp_keepalive": True, }注意:指定空字典
{}会覆盖连接配置中的botocore.config.Config设置。
启动 Neptune 数据库集群
使用NeptuneStartDbClusterOperator启动已存在的 Neptune 集群。运算符支持可延迟模式:传入deferrable=True即异步等待集群启动完成,此模式要求安装aiobotocore模块。
start_cluster = NeptuneStartDbClusterOperator(task_id="start_task", db_cluster_id=cluster_id)该示例来自系统测试 example_neptune.py。
参数说明
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
db_cluster_id | str | 必填 | 要启动的 Neptune 集群标识符 |
wait_for_completion | bool | True | 是否等待集群启动完成 |
deferrable | bool | 配置项operators.default_deferrable(默认False) | 若为True,异步等待集群启动(隐含等待完成),需要aiobotocore |
waiter_delay | int | 30 | 状态检查间隔(秒) |
waiter_max_attempts | int | 60 | 最大检查次数 |
执行返回值为字典{"db_cluster_id": cluster_id}。
执行流程与状态机
从源码 neptune.py 可见,启动流程如下:
- 通过
NeptuneHook.get_cluster_status查询集群当前状态; - 若状态在
AVAILABLE_STATES(available)中,直接返回,不重复启动; - 若状态在
ERROR_STATES中(cloning-failed、inaccessible-encryption-credentials、inaccessible-encryption-credentials-recoverable、migration-failed),抛出AirflowException,因为错误状态下无法启动; - 调用
conn.start_db_cluster(DBClusterIdentifier=...); - 若抛出可等待的
ClientError(如InvalidDBInstanceState、InvalidClusterState、InvalidDBClusterStateFault),通过handle_waitable_exception等待集群/实例可用后重试; - 若
deferrable=True,委托NeptuneClusterAvailableTrigger异步等待; - 否则若
wait_for_completion=True,调用hook.wait_for_cluster_availability同步轮询。
底层等待逻辑由 NeptuneHook 实现,其核心是使用 waitercluster_available,配置定义在 neptune.json:
success:DBClusters[0].Status == "available"failure:状态为deleting、inaccessible-encryption-credentials、inaccessible-encryption-credentials-recoverable、migration-failedretry:状态为stopped(继续轮询)
集群与实例状态联动:集群与其实例必须都处于有效状态才能发送启动请求。当遇到InvalidDBInstanceState错误时,运算符会先等待db_instance_availablewaiter(通过NeptuneClusterInstancesAvailableTrigger或hook.wait_for_cluster_instance_availability),再重试启动。单元测试 test_neptune.py 验证了该行为。
停止 Neptune 数据库集群
使用NeptuneStopDbClusterOperator停止运行中的 Neptune 集群。同样支持deferrable=True可延迟模式(需安装aiobotocore)。
stop_cluster = NeptuneStopDbClusterOperator(task_id="stop_task", db_cluster_id=cluster_id)该示例同样来自 example_neptune.py。
参数与启动运算符完全一致(db_cluster_id、wait_for_completion、deferrable、waiter_delay、waiter_max_attempts),返回值同样为{"db_cluster_id": cluster_id}。
执行流程与状态机
停止流程与启动对称(见 neptune.py):
- 查询集群状态;
- 若状态在
STOPPED_STATES(stopped)中,直接返回,不重复停止; - 若状态在
ERROR_STATES中,抛出AirflowException; - 调用
conn.stop_db_cluster(DBClusterIdentifier=...); - 遇到可等待的
ClientError时同样进入等待重试逻辑; deferrable=True时委托NeptuneClusterStoppedTrigger;- 否则若
wait_for_completion=True,调用hook.wait_for_cluster_stopped。
停止等待使用的 waitercluster_stopped(定义见 neptune.json):
success:DBClusters[0].Status == "stopped"failure:状态为deleting、inaccessible-encryption-credentials、inaccessible-encryption-credentials-recoverable、migration-failed
可延迟(Deferrable)执行模式
两个运算符都支持deferrable=True,该模式的核心价值是:任务在触发异步等待后立即释放 worker 槽位,由 Trigger 在后台轮询状态,状态满足后再唤醒任务继续执行,从而显著降低资源占用。此模式需要安装aiobotocore:
start_cluster = NeptuneStartDbClusterOperator( task_id="start_task", db_cluster_id=cluster_id, deferrable=True, waiter_delay=30, waiter_max_attempts=60, )相关 Trigger 定义在 triggers/neptune.py:
| Trigger | 等待目标 | 轮询的 waiter |
|---|---|---|
NeptuneClusterAvailableTrigger | 集群可用 | cluster_available |
NeptuneClusterStoppedTrigger | 集群停止 | cluster_stopped |
NeptuneClusterInstancesAvailableTrigger | 集群实例可用 | db_instance_available |
Trigger 通过AwsBaseWaiterTrigger复用自定义 waiter,并序列化db_cluster_id、aws_conn_id、region_name、waiter_delay、waiter_max_attempts等字段。单元测试 test_neptune.py 验证了 operator 的配置会正确传递到 Trigger(包括waiter_delay与waiter_max_attempts),确保异步等待使用与同步模式一致的轮询参数。
deferrable的默认值取自 Airflow 配置项operators.default_deferrable(源码见 neptune.py),因此你也可以在airflow.cfg中全局开启默认可延迟行为。
在 DAG 中组合使用
参考系统测试 example_neptune.py 的完整编排(创建 → 启动 → 停止 → 删除):
from datetime import datetime from airflow import DAG from airflow.providers.amazon.aws.operators.neptune import ( NeptuneStartDbClusterOperator, NeptuneStopDbClusterOperator, ) with DAG( dag_id="example_neptune", start_date=datetime(2021, 1, 1), schedule="@once", catchup=False, ) as dag: # [START howto_operator_start_neptune_cluster] start_cluster = NeptuneStartDbClusterOperator(task_id="start_task", db_cluster_id=cluster_id) # [END howto_operator_start_neptune_cluster] # [START howto_operator_stop_neptune_cluster] stop_cluster = NeptuneStopDbClusterOperator(task_id="stop_task", db_cluster_id=cluster_id) # [END howto_operator_stop_neptune_cluster]在实际生产 DAG 中,可以配合其他任务实现"定时开机/关机"的省钱策略,例如在非工作时间停止集群、工作时段自动启动。
参考
- 运算符源码与类文档:operators/neptune.py
- Hook 实现(状态常量与等待方法):hooks/neptune.py
- Trigger 实现:triggers/neptune.py
- Waiter 定义:waiters/neptune.json
- 系统测试示例:example_neptune.py
- 单元测试:test_neptune.py
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考