Apache Airflow 与 Amazon Neptune:使用 Start / Stop 运算符管理图数据库集群生命周期
2026/9/13 6:32:38 网站建设 项目流程

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 提供了两个专门运算符NeptuneStartDbClusterOperatorNeptuneStopDbClusterOperator,让你能够以可编程方式启动、停止 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_idAWS 连接 ID。若设为None,则使用默认的 boto3 行为(不进行连接查询);否则使用连接中存储的凭据aws_default
region_nameAWS 区域名。若为None,使用 AWS 连接 Extra 参数中的 region_name;否则覆盖连接值None
verify是否校验 SSL 证书。False表示不校验;也可指定 CA 证书 bundle 文件路径。若为None,使用连接 Extra 参数中的 verifyNone
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_idstr必填要启动的 Neptune 集群标识符
wait_for_completionboolTrue是否等待集群启动完成
deferrablebool配置项operators.default_deferrable(默认False若为True,异步等待集群启动(隐含等待完成),需要aiobotocore
waiter_delayint30状态检查间隔(秒)
waiter_max_attemptsint60最大检查次数

执行返回值为字典{"db_cluster_id": cluster_id}

执行流程与状态机

从源码 neptune.py 可见,启动流程如下:

  1. 通过NeptuneHook.get_cluster_status查询集群当前状态;
  2. 若状态在AVAILABLE_STATESavailable)中,直接返回,不重复启动;
  3. 若状态在ERROR_STATES中(cloning-failedinaccessible-encryption-credentialsinaccessible-encryption-credentials-recoverablemigration-failed),抛出AirflowException,因为错误状态下无法启动;
  4. 调用conn.start_db_cluster(DBClusterIdentifier=...)
  5. 若抛出可等待的ClientError(如InvalidDBInstanceStateInvalidClusterStateInvalidDBClusterStateFault),通过handle_waitable_exception等待集群/实例可用后重试;
  6. deferrable=True,委托NeptuneClusterAvailableTrigger异步等待;
  7. 否则若wait_for_completion=True,调用hook.wait_for_cluster_availability同步轮询。

底层等待逻辑由 NeptuneHook 实现,其核心是使用 waitercluster_available,配置定义在 neptune.json:

  • successDBClusters[0].Status == "available"
  • failure:状态为deletinginaccessible-encryption-credentialsinaccessible-encryption-credentials-recoverablemigration-failed
  • retry:状态为stopped(继续轮询)

集群与实例状态联动:集群与其实例必须都处于有效状态才能发送启动请求。当遇到InvalidDBInstanceState错误时,运算符会先等待db_instance_availablewaiter(通过NeptuneClusterInstancesAvailableTriggerhook.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_idwait_for_completiondeferrablewaiter_delaywaiter_max_attempts),返回值同样为{"db_cluster_id": cluster_id}

执行流程与状态机

停止流程与启动对称(见 neptune.py):

  1. 查询集群状态;
  2. 若状态在STOPPED_STATESstopped)中,直接返回,不重复停止;
  3. 若状态在ERROR_STATES中,抛出AirflowException
  4. 调用conn.stop_db_cluster(DBClusterIdentifier=...)
  5. 遇到可等待的ClientError时同样进入等待重试逻辑;
  6. deferrable=True时委托NeptuneClusterStoppedTrigger
  7. 否则若wait_for_completion=True,调用hook.wait_for_cluster_stopped

停止等待使用的 waitercluster_stopped(定义见 neptune.json):

  • successDBClusters[0].Status == "stopped"
  • failure:状态为deletinginaccessible-encryption-credentialsinaccessible-encryption-credentials-recoverablemigration-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_idaws_conn_idregion_namewaiter_delaywaiter_max_attempts等字段。单元测试 test_neptune.py 验证了 operator 的配置会正确传递到 Trigger(包括waiter_delaywaiter_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),仅供参考

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

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

立即咨询