ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Airflow 与 Amazon Neptune:使用 Start / Stop 运算符管理图数据库集群生命周期

Apache Airflow 与 Amazon Neptune:使用 Start / Stop 运算符管理图数据库集群生命周期 Apache Airflow 与 Amazon Neptune使用 Start / Stop 运算符管理图数据库集群生命周期【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon Neptune 是无服务器图数据库服务专为高性能、高可扩展性与高可用性设计内置安全能力、持续备份以及与其它 AWS 服务的集成。Apache Airflow 的 amazon provider 提供了两个专门运算符NeptuneStartDbClusterOperator与NeptuneStopDbClusterOperator让你能够以可编程方式启动、停止 Neptune 数据库集群并自动等待集群达到目标状态。本文基于官方文档 neptune.rst结合源码、Hook、Trigger、Waiter 配置与测试代码全面讲解这两个运算符的使用方法、参数语义、可延迟deferrable执行模式以及底层实现原理帮助你直接在 DAG 中安全、高效地管理 Neptune 集群生命周期。前置准备1. 安装 amazon providerpip 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_defaultregion_nameAWS 区域名。若为None使用 AWS 连接 Extra 参数中的 region_name否则覆盖连接值Noneverify是否校验 SSL 证书。False表示不校验也可指定 CA 证书 bundle 文件路径。若为None使用连接 Extra 参数中的 verifyNonebotocore_config用于构造botocore.config.Config的字典可配置重试策略、超时等Nonebotocore_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 集群。运算符支持可延迟模式传入deferrableTrue即异步等待集群启动完成此模式要求安装aiobotocore模块。start_cluster NeptuneStartDbClusterOperator(task_idstart_task, db_cluster_idcluster_id)该示例来自系统测试 example_neptune.py。参数说明参数类型默认值说明db_cluster_idstr必填要启动的 Neptune 集群标识符wait_for_completionboolTrue是否等待集群启动完成deferrablebool配置项operators.default_deferrable默认False若为True异步等待集群启动隐含等待完成需要aiobotocorewaiter_delayint30状态检查间隔秒waiter_max_attemptsint60最大检查次数执行返回值为字典{db_cluster_id: cluster_id}。执行流程与状态机从源码 neptune.py 可见启动流程如下通过NeptuneHook.get_cluster_status查询集群当前状态若状态在AVAILABLE_STATESavailable中直接返回不重复启动若状态在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等待集群/实例可用后重试若deferrableTrue委托NeptuneClusterAvailableTrigger异步等待否则若wait_for_completionTrue调用hook.wait_for_cluster_availability同步轮询。底层等待逻辑由 NeptuneHook 实现其核心是使用 waitercluster_available配置定义在 neptune.jsonsuccessDBClusters[0].Status availablefailure状态为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 集群。同样支持deferrableTrue可延迟模式需安装aiobotocore。stop_cluster NeptuneStopDbClusterOperator(task_idstop_task, db_cluster_idcluster_id)该示例同样来自 example_neptune.py。参数与启动运算符完全一致db_cluster_id、wait_for_completion、deferrable、waiter_delay、waiter_max_attempts返回值同样为{db_cluster_id: cluster_id}。执行流程与状态机停止流程与启动对称见 neptune.py查询集群状态若状态在STOPPED_STATESstopped中直接返回不重复停止若状态在ERROR_STATES中抛出AirflowException调用conn.stop_db_cluster(DBClusterIdentifier...)遇到可等待的ClientError时同样进入等待重试逻辑deferrableTrue时委托NeptuneClusterStoppedTrigger否则若wait_for_completionTrue调用hook.wait_for_cluster_stopped。停止等待使用的 waitercluster_stopped定义见 neptune.jsonsuccessDBClusters[0].Status stoppedfailure状态为deleting、inaccessible-encryption-credentials、inaccessible-encryption-credentials-recoverable、migration-failed可延迟Deferrable执行模式两个运算符都支持deferrableTrue该模式的核心价值是任务在触发异步等待后立即释放 worker 槽位由 Trigger 在后台轮询状态状态满足后再唤醒任务继续执行从而显著降低资源占用。此模式需要安装aiobotocorestart_cluster NeptuneStartDbClusterOperator( task_idstart_task, db_cluster_idcluster_id, deferrableTrue, waiter_delay30, waiter_max_attempts60, )相关 Trigger 定义在 triggers/neptune.pyTrigger等待目标轮询的 waiterNeptuneClusterAvailableTrigger集群可用cluster_availableNeptuneClusterStoppedTrigger集群停止cluster_stoppedNeptuneClusterInstancesAvailableTrigger集群实例可用db_instance_availableTrigger 通过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_idexample_neptune, start_datedatetime(2021, 1, 1), scheduleonce, catchupFalse, ) as dag: # [START howto_operator_start_neptune_cluster] start_cluster NeptuneStartDbClusterOperator(task_idstart_task, db_cluster_idcluster_id) # [END howto_operator_start_neptune_cluster] # [START howto_operator_stop_neptune_cluster] stop_cluster NeptuneStopDbClusterOperator(task_idstop_task, db_cluster_idcluster_id) # [END howto_operator_stop_neptune_cluster]在实际生产 DAG 中可以配合其他任务实现定时开机/关机的省钱策略例如在非工作时间停止集群、工作时段自动启动。参考运算符源码与类文档operators/neptune.pyHook 实现状态常量与等待方法hooks/neptune.pyTrigger 实现triggers/neptune.pyWaiter 定义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),仅供参考
返回列表