
Apache Airflow 集成 Amazon OpenSearch ServerlessCollection 创建与状态等待实战指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon OpenSearch Serverless 是 Amazon OpenSearch Service 的按需、自动扩缩容形态Collection集合会随应用负载自动调整计算容量与需要人工规划容量的 Provisioned 域形成鲜明对比。本文基于 Apache Airflow Amazon Provider 官方文档 opensearchserverless.rst 及其配套源码系统讲解如何在 Airflow DAG 中使用OpenSearchServerlessCreateCollectionOperator创建 Collection、用OpenSearchServerlessCollectionActiveSensor等待其进入 ACTIVE 状态并深入剖析背后的 Hook、Trigger 与自定义 Waiter 实现让你能直接上手编排 Serverless 搜索与向量检索基础设施。为什么用 OpenSearch Serverless 而非 Provisioned 域Amazon OpenSearch Serverless 是 Amazon OpenSearch Service 的按需、自动扩缩容配置。一个 OpenSearch ServerlessCollection本质上是一个根据应用需求自动扩缩计算能力的 OpenSearch 集群这与需要人工管理容量的 Provisioned OpenSearch 域形成鲜明对比。对于 RAG检索增强生成知识库、日志检索、向量相似度搜索等负载波动明显的场景Serverless 形态可以避免容量预估与闲置成本问题。从代码组织看Amazon Provider 为 OpenSearch Serverless 提供了完整的能力栈包括Hookopensearch_serverless.pyOperatoropensearch_serverless.pySensoropensearch_serverless.pyTriggeropensearch_serverless.py自定义 Waiter 定义opensearchserverless.json前置任务Prerequisite Tasks要使用本文涉及的 Operator 与 Sensor需要完成三件事创建 AWS 资源通过 AWS Console 或 AWS CLI 准备好运行所需的 AWS 资源。安装 Python 依赖通过 pip 安装 Amazon Providerpip install apache-airflow[amazon]更完整的安装指引可参考 Airflow 的安装文档见仓库 airflow-core/docs/installation 目录。配置 AWS Connection在 Airflow 中配置 Amazon Web Services 连接详细步骤见 Setup Connection。这些前置任务在官方文档中通过 prerequisite_tasks.rst 片段被所有 AWS 相关指南共享。通用参数Generic Parameters所有 OpenSearch Serverless 相关的 Operator、Sensor 都继承自AwsBaseOperator/AwsBaseSensor因此支持一组通用的 AWS 参数定义于 generic_parameters.rst参数说明默认值aws_conn_id引用的 Airflow AWS Connection ID设为None时使用默认 boto3 行为不查 Connection否则使用 Connection 中存储的凭证aws_defaultregion_nameAWS 区域名为None或省略时使用 AWS Connection Extra 参数中的region_name否则覆盖之Noneverify是否校验 SSL 证书。False表示不校验也可传入 CA 证书 bundle 文件路径如path/to/cert/bundle.pem以使用与 botocore 不同的证书链为None或省略时使用 Connection 中的verify配置Nonebotocore_config用于构造botocore.config.Config的字典可配置重试策略、超时、限流规避throttling等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, }需要注意显式传入空字典{}会覆盖Connection 中的botocore.config.Config配置。使用 Sensor 等待 Collection 变为 ACTIVEOpenSearchServerlessCollectionActiveSensor 简介创建 Collection 是异步操作创建完成后需要等待其进入可用状态才能进行索引写入、知识库挂载等后续步骤。官方文档推荐使用OpenSearchServerlessCollectionActiveSensor位于 sensors/opensearch_serverless.py来轮询 Collection 状态直至达到终态。官方系统测试中给出的用法如下见 example_bedrock_retrieve_and_generate.pyawait_collection OpenSearchServerlessCollectionActiveSensor( task_idawait_collection, collection_namevector_store_name, )核心参数参数说明默认值collection_idCollection 的 ID与collection_name二选一不能同时提供Nonecollection_nameCollection 的名称与collection_id二选一不能同时提供Nonepoke_interval轮询间隔秒10max_retries最大重试次数超过后返回当前状态60deferrable是否以可延迟deferrable模式运行开启时需要安装aiobotocore。默认取自配置文件operators.default_deferrableFalse状态机与轮询逻辑从源码看Sensor 内部定义了明确的状态集合INTERMEDIATE_STATES (CREATING,) FAILURE_STATES (DELETING, FAILED,) SUCCESS_STATES (ACTIVE,)每次poke()调用会通过hook.conn.batch_get_collection()查询状态状态属于FAILURE_STATESDELETING/FAILED→ 抛出AirflowException任务失败状态属于INTERMEDIATE_STATESCREATING→ 返回False继续轮询否则即ACTIVE→ 返回True任务成功。可延迟执行Deferrable模式当deferrableTrue时Sensor 不会占用 worker 轮询而是通过defer()将自身挂起并交给OpenSearchServerlessCollectionActiveTrigger见 triggers/opensearch_serverless.py异步等待。该 Trigger 继承自AwsBaseWaiterTrigger复用仓库自定义的collection_availableWaiter定义见 opensearchserverless.json{ version: 2, waiters: { collection_available: { operation: BatchGetCollection, delay: 10, maxAttempts: 120, acceptors: [ {matcher: path, argument: collectionDetails[0].status, expected: ACTIVE, state: success}, {matcher: path, argument: collectionDetails[0].status, expected: DELETING, state: failure}, {matcher: path, argument: collectionDetails[0].status, expected: CREATING, state: retry}, {matcher: path, argument: collectionDetails[0].status, expected: FAILED, state: failure} ] } } }该 Waiter 基于BatchGetCollection操作以 10 秒为间隔、最多 120 次尝试合计约 20 分钟通过路径匹配collectionDetails[0].status字段判定成功ACTIVE、失败DELETING/FAILED或重试CREATING。Trigger 模式下waiter_delay默认 60 秒、waiter_max_attempts默认 20 次。使用 Operator 创建 CollectionOpenSearchServerlessCreateCollectionOperator 简介官方文档推荐使用OpenSearchServerlessCreateCollectionOperator位于 operators/opensearch_serverless.py创建 Collection其官方系统测试示例见 example_opensearch_serverless.pycreate_collection OpenSearchServerlessCreateCollectionOperator( task_idcreate_collection, collection_namecollection_name, collection_typeSEARCH, )核心参数参数说明默认值collection_nameCollection 名称支持模板渲染必填collection_typeCollection 类型SEARCH、TIMESERIES、VECTORSEARCH支持模板渲染SEARCHdescription可选描述支持模板渲染Nonestandby_replicas是否启用备用副本ENABLED或DISABLEDNonetags可选的标签字典列表如[{key: env, value: prod}]Noneif_exists当 Collection 已存在时的行为fail抛出错误skip记录日志并返回skip其中collection_name、collection_type、description被声明为模板字段template_fields支持 Jinja 模板动态渲染。执行逻辑与返回值从源码看execute()的核心流程如下使用prune_dict剔除值为None的参数构造create_collection调用参数调用self.hook.conn.create_collection(**kwargs)发起创建从响应response[createCollectionDetail][id]中取出新 Collection 的 ID若抛出ClientError且错误码为ConflictException当if_exists skip时调用batch_get_collection(names[...])查询已存在的 Collection 并返回其 ID否则重新抛出异常。最终返回 Collection ID字符串该返回值可通过 XCom 被下游任务引用。Hook 层封装Operator 与 Sensor 底层的OpenSearchServerlessHook见 hooks/opensearch_serverless.py是对boto3.client(opensearchserverless)的薄封装class OpenSearchServerlessHook(AwsBaseHook): client_type opensearchserverless def __init__(self, *args, **kwargs) - None: kwargs[client_type] self.client_type super().__init__(*args, **kwargs)它继承AwsBaseHook自动复用aws_conn_id、region_name、verify、botocore_config等通用参数与凭证解析逻辑任务代码中通过self.hook.conn即可直接访问 boto3 客户端。完整实战创建 Collection 并等待就绪将 Operator 与 Sensor 组合起来即可构成一个典型的创建基础设施 → 等待就绪DAG。参考官方系统测试 example_opensearch_serverless.py 的完整流程from datetime import datetime from airflow.providers.amazon.aws.operators.opensearch_serverless import ( OpenSearchServerlessCreateCollectionOperator, ) from airflow.providers.amazon.aws.sensors.opensearch_serverless import ( OpenSearchServerlessCollectionActiveSensor, ) from airflow.sdk import DAG, chain with DAG( dag_idexample_opensearch_serverless, scheduleNone, start_datedatetime(2024, 1, 1), catchupFalse, ) as dag: collection_name my-serverless-collection create_collection OpenSearchServerlessCreateCollectionOperator( task_idcreate_collection, collection_namecollection_name, collection_typeSEARCH, ) wait_for_collection OpenSearchServerlessCollectionActiveSensor( task_idwait_for_collection, collection_namecollection_name, poke_interval30, timeout600, ) chain(create_collection, wait_for_collection)系统测试中还展示了与 Bedrock 知识库集成的真实场景example_bedrock_retrieve_and_generate.py先用create_collection创建向量存储 Collection再用OpenSearchServerlessCollectionActiveSensor等待其 ACTIVE随后才创建 Bedrock Knowledge Base 并挂载该 Collection 作为向量索引——这印证了先建 Collection、再等就绪、后挂载使用的标准编排顺序。测试佐证仓库单元测试覆盖了 Operator 与 Sensor 的完整行为见 tests/unit/amazon/aws/operators/test_opensearch_serverless.py 与 tests/unit/amazon/aws/sensors/test_opensearch_serverless.py包括ConflictException跳过逻辑、ID 与名称互斥校验、可延迟模式下的 Trigger 委托等关键路径可作为理解行为边界的参考。参考AWS boto3 官方文档Amazon OpenSearch Service 与 Amazon OpenSearch Serverless 服务客户端参考原文档参考章节所列可查阅对应服务 API 细节官方指南原文operators/opensearchserverless.rstAWS Connection 配置docs/connections/aws.rst系统测试示例example_opensearch_serverless.py【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考