Apache Airflow 中的 SageMaker Unified Studio 算子:编排 Notebook、Querybook 与 Visual ETL 作业
Apache Airflow 中的 SageMaker Unified Studio 算子编排 Notebook、Querybook 与 Visual ETL 作业【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon SageMaker Unified Studio 是 AWS 提供的一种统一开发体验它把 AWS 的数据、分析、人工智能AI与机器学习ML服务汇聚到同一个界面中让团队能够在一个地方构建、部署、执行并监控端到端的工作流从而促进跨团队协作与敏捷开发。Apache Airflow 的 Amazon 提供商包apache-airflow-providers-amazon为此提供了一组专门用于在 SageMaker Unified Studio 中运行各类工件artifact的算子。本文以 sagemakerunifiedstudio.rst 为主线结合仓库中的算子、Hook、Sensor、Trigger 源码与系统测试 DAG系统讲解如何用 Airflow 驱动 SageMaker Unified Studio 中的 Notebook、Querybook、Visual ETL 作业与 Unified Studio Notebook读完即可在自己的 DAG 中落地使用。前置条件使用这些算子的准备工作在编写 DAG 之前需要先在 AWS 侧完成以下准备对应 官方入门文档 的步骤创建 SageMaker Unified Studio 域Domain与项目Project算子中的domain_id/domain_identifier与project_id/owning_project_identifier都指向这两个资源。如果域是 IdCIAM Identity Center域进入 Compute Workflow environments 选项卡点击 Create 创建一个新的MWAA 环境。MWAA 环境用于为算子提供运行时解析所需的上下文。准备工件Artifact创建一个 Jupyter notebook、querybook、Visual ETL 作业或 SageMaker Unified Studio notebook并将其保存到你的项目中。两个算子正是围绕这些工件类型设计的。从系统测试的前置说明见 example_sagemaker_unified_studio.py还可以看到测试账号还要求预先配置一个 IAM IDC 组织、一个带默认 VPC 与角色的 SageMaker Unified Studio Domain、Domain 内的项目以及放在项目 S3 路径下的test_notebook.ipynb。算子选型两条执行路径的区别Airflow 为 SageMaker Unified Studio 提供了两个不同的算子它们的执行机制和工件标识方式完全不同应根据使用场景选择算子执行工件底层机制工件标识方式SageMakerNotebookOperatorJupyter notebook、querybook、Visual ETL 作业依赖sagemaker_studioPython 库项目内的相对文件路径如test_notebook.ipynbSageMakerUnifiedStudioNotebookOperatorSageMaker Unified Studio notebook通过 DataZoneStartNotebookRunAPIboto3datazoneclientNotebook ID如nb-1234567890加 Domain ID、Project ID从源码看两者连 Hook 体系也是完全分开的前者走 SageMakerNotebookHook包装sagemaker_studioSDK 的 Execution API后者走 SageMakerUnifiedStudioNotebookHook包装 DataZone NotebookRun API。下面分别展开。SageMakerNotebookOperator运行 Jupyter Notebook、Querybook 与 Visual ETL 作业SageMakerNotebookOperator用于执行 Jupyter notebooks、querybooks 和 Visual ETL 作业底层依赖sagemaker_studioPython 库。工件通过项目内的相对文件路径标识例如test_notebook.ipynb也可以是.sqlnb的 querybook 或.vetl的 Visual ETL 文件。参数详解与默认值该算子的完整签名定义在 sagemaker_unified_studio.py核心参数如下input_config必填输入文件配置。input_path必须是相对路径算子在执行时会自动将其解析为 IDE 中用户主目录上下文下的绝对路径input_params为传给工件的参数字典。示例{input_path: folder/input/notebook.ipynb, input_params: {key: value}}。源码会在执行前校验input_config与input_path的存在性缺失直接抛AirflowException。domain_id可选SageMaker Unified Studio 的 Domain ID。不传时 SDK 会尝试从环境变量解析。project_id可选项目 ID同样支持环境变量解析。domain_region可选Domain 所在的 AWS 区域不传时使用默认 AWS 区域。output_config可选输出格式配置例如{output_formats: [NOTEBOOK]}默认值即{output_formats: [NOTEBOOK]}。compute可选远程执行的计算资源配置远程计算上执行时为必填。示例{ instance_type: ml.c5.xlarge, image_details: { image_name: sagemaker-distribution-prod, image_version: 3, ecr_uri: 123456123456.dkr.ecr.us-west-2.amazonaws.com/ImageName:latest, }, }termination_condition可选远程执行的终止条件例如{MaxRuntimeInSeconds: 3600}。tags可选附加到远程执行运行的标签例如{md_analytics: logs}。wait_for_completion默认True是否等待 Notebook 执行完成为False时发起执行后立即返回。waiter_delay默认10轮询执行状态的时间间隔秒。waiter_max_attempts默认1440超过该次数仍未完成则任务以失败结束。deferrable默认False为True时算子异步等待任务完成隐式要求等待完成需要安装aiobotocore模块。此外domain_id、project_id、domain_region、input_config、output_config、compute、termination_condition、tags均被声明为模板字段template_fields支持 Jinja 模板渲染。方式一显式传入 Domain / Project / Region 参数推荐这是文档推荐的第一种用法把domain_id、project_id、domain_region直接作为算子参数传入无需任何环境变量——SDK 会根据这些参数自行解析 S3 路径与区域。注意该方式要求sagemaker-studio1.0.25。示例来自 example_sagemaker_unified_studio.py 的howto_operator_sagemaker_unified_studio_notebook_explicit_params段落# Run notebook with domain_id/project_id/domain_region passed explicitly as operator parameters. # No environment variables needed — the SDK resolves the S3 path and region from these params. # Requires sagemaker-studio1.0.25. run_notebook_explicit_params SageMakerNotebookOperator( task_idrun-notebook-explicit, domain_iddomain_id, project_idproject_id, domain_regionregion_name, input_config{input_path: notebook_path, input_params: {}}, output_config{output_formats: [NOTEBOOK]}, # optional compute{ instance_type: ml.m5.large, volume_size_in_gb: 30, }, # optional termination_condition{max_runtime_in_seconds: 600}, # optional tags{}, # optional wait_for_completionTrue, # optional waiter_delay5, # optional deferrableFalse, # optional )方式二基于 MWAA 风格环境变量的解析路径传统方式第二种用法不传 Domain/Project 参数而是依赖一组AIRFLOW__WORKFLOWS__*前缀的环境变量来解析上下文。这正对应前置条件中提到的MWAA 环境在 SageMaker Unified Studio 中创建的 MWAA 环境会自动注入这些变量。系统测试中模拟该环境的变量集合见 example_sagemaker_unified_studio.py 的get_mwaa_environment_paramsAIRFLOW__WORKFLOWS__DATAZONE_DOMAIN_ID # Domain ID AIRFLOW__WORKFLOWS__DATAZONE_PROJECT_ID # Project ID AIRFLOW__WORKFLOWS__DATAZONE_ENVIRONMENT_ID # 环境 ID AIRFLOW__WORKFLOWS__DATAZONE_SCOPE_NAME # 例如 dev AIRFLOW__WORKFLOWS__DATAZONE_STAGE # 例如 prod AIRFLOW__WORKFLOWS__DATAZONE_ENDPOINT # 例如 https://datazone.{region}.api.aws AIRFLOW__WORKFLOWS__PROJECT_S3_PATH # 项目 S3 路径 AIRFLOW__WORKFLOWS__DATAZONE_DOMAIN_REGION # Domain 所在区域对应的 DAG 代码howto_operator_sagemaker_unified_studio_notebook段落# Run notebook using the legacy env-var-based resolution path (MWAA-style). run_notebook SageMakerNotebookOperator( task_idrun-notebook, input_config{input_path: notebook_path, input_params: {}}, output_config{output_formats: [NOTEBOOK]}, # optional compute{ instance_type: ml.m5.large, volume_size_in_gb: 30, }, # optional termination_condition{max_runtime_in_seconds: 600}, # optional tags{}, # optional wait_for_completionTrue, # optional waiter_delay5, # optional deferrableFalse, # optional executor_config{ # optional overrides: { containerOverrides: [ { environment: [ {name: key, value: value} for key, value in mock_mwaa_environment_params.items() ], name: ECSExecutorContainer, # Necessary parameter } ] } }, )系统测试刻意让“显式参数”任务在设置环境变量之前运行以证明显式参数在没有任何 MWAA 风格环境变量的情况下也能独立工作而旧版 SDK1.0.25无法从显式参数解析区域因此只能用方式二。底层执行原理Hook 与状态轮询从源码看SageMakerNotebookOperator.execute()的执行链路是通过notebook_execution_hook懒加载的SageMakerNotebookHook调用start_notebook_execution()拿到execution_id。若deferrableTrue将任务委托给 SageMakerNotebookJobTrigger 异步等待使用asyncio.to_thread轮询get_notebook_execution免去run_in_executor对 CI 更友好否则若wait_for_completionTrue同步轮询直到终态。Hook 在构造时通过ClientConfig组装配置并把local、domain_identifier、project_identifier、datazone_domain_region写入overrides[execution]其中local来自is_local_runner()即检测环境变量WORKFLOWS_ENV是否等于Local见 utils/sagemaker_unified_studio.py。启动执行时它会组装execution_typeNOTEBOOK的请求把input_config规范化为{notebook_config: {input_path: ..., input_parameters: ...}}把output_config规范化为{notebook_config: {output_formats: [...]}}。同步轮询wait_for_execution_completion()的核心逻辑IN_PROGRESS、STOPPING视为进行中按waiter_delay间隔继续轮询COMPLETED视为成功并返回{Status: ..., ExecutionId: ...}其余状态视为失败抛出AirflowException附上 API 返回的error_message超过waiter_max_attempts仍未完成则判定超时失败。轮询过程中如果响应里带有files输出文件列表或s3_pathHook 还会把它们推送到 XComfiles以{display_name}.{file_format}为 key、file_path为 values3_path则直接以s3_path为 key。SageMakerUnifiedStudioNotebookOperator通过 DataZone API 执行 Unified Studio NotebookSageMakerUnifiedStudioNotebookOperator用于执行SageMaker Unified Studio 的 Notebook 工件它不再依赖sagemaker_studioSDK而是直接通过 DataZone 的StartNotebookRunAPI 发起无头headless执行。Notebook 由Notebook ID如nb-1234567890连同其所在的 Domain ID 与 Project ID 共同标识。参数详解与默认值定义见 sagemaker_unified_studio_notebook.py继承自AwsBaseOperator走标准 AWS 连接/凭证体系notebook_identifier必填要执行的 Notebook ID。domain_identifier必填Notebook 所在的 SageMaker Unified StudioDataZoneDomain ID。owning_project_identifier必填Notebook 所属项目 ID。client_token可选幂等令牌不传时 Hook 自动生成uuid4()。notebook_parameters可选传给 Notebook 的参数字典支持 Jinja 模板模板字段。compute_configuration可选计算配置例如{instanceType: sc.m5.large}。timeout_configuration可选超时设置例如{runTimeoutInMinutes: 1440}轮询的最大次数由runTimeoutInMinutes * 60 / waiter_delay推导不传时默认 12 小时。wait_for_completion默认True等待运行结束False时发起运行后立即返回{notebook_run_id: ...}。waiter_delay默认10轮询间隔秒必须为正整数。deferrable默认False为True时交给 Trigger 异步等待释放 worker 槽位。endpoint_url可选DataZone API 的自定义 endpoint。基础用法示例来自 example_sagemaker_unified_studio_notebook.py 的howto_operator_sagemaker_unified_studio_notebook段落client_token fidempotency-token-{int(time.time())} run_notebook SageMakerUnifiedStudioNotebookOperator( task_idnotebook-task, aws_conn_idDATAZONE_CONN_ID, notebook_identifiernotebook_id, domain_identifierdomain_id, owning_project_identifierproject_id, client_tokenclient_token, # optional notebook_parameters{ param1: value1, param2: value2, }, # optional compute_configuration{instanceType: sc.m5.large}, # optional timeout_configuration{runTimeoutInMinutes: 1440}, # optional wait_for_completionTrue, # optional waiter_delay30, # optional deferrableFalse, # optional )系统测试还展示了配套的 AWS 连接配置方式通过一个conn_typeaws、extra中携带{role_arn: ..., assume_role_method: assume_role}的连接aws_datazone_notebook来让算子代入 DataZone 环境角色。Notebook 输出变量与 XCom 传递Notebook A → Notebook BSageMakerUnifiedStudioNotebookOperator的一个亮点是Notebook 可以产出输出变量运行完成后自动推送到 XCom下游任务可以通过notebook_parameters中的 Jinja 模板消费这些输出。示例howto_operator_sagemaker_unified_studio_notebook_pass_outputs段落# Notebook A produces outputs (e.g., name, age) that are pushed to xcom. # Notebook B consumes those outputs via Jinja templating in notebook_parameters. run_notebook_a SageMakerUnifiedStudioNotebookOperator( task_idnotebook-a-task, aws_conn_idDATAZONE_CONN_ID, notebook_identifiernotebook_id, domain_identifierdomain_id, owning_project_identifierproject_id, wait_for_completionTrue, deferrableFalse, ) run_notebook_b SageMakerUnifiedStudioNotebookOperator( task_idnotebook-b-task, aws_conn_idDATAZONE_CONN_ID, notebook_identifiernotebook_b_id, domain_identifierdomain_id, owning_project_identifierproject_id, notebook_parameters{ employee_name: {{ task_instance.xcom_pull(task_idsnotebook-a-task, keyNOTEBOOK_OUTPUT.name) }}, employee_age: {{ task_instance.xcom_pull(task_idsnotebook-a-task, keyNOTEBOOK_OUTPUT.age) }}, }, wait_for_completionTrue, deferrableFalse, )这里 Notebook A 产出name与age两个输出Notebook B 通过task_instance.xcom_pull(task_idsnotebook-a-task, keyNOTEBOOK_OUTPUT.name)之类的 Jinja 表达式接收。实现上见 sagemaker_unified_studio_notebook.py 与算子源码的_push_notebook_outputs运行结束后SDK 会把输出变量以 JSON 形式写入项目 S3 桶中的固定位置.sys/notebooks/{notebook_identifier}/runs/{notebook_run_id}/notebook_outputs.jsonIDC 域会带上项目前缀Hook 通过get_project_s3_path()解析项目的 S3 桶与前缀——它调用 DataZone API 找到项目默认的Tooling找不到则回退ToolingLite环境读取其s3BucketPathprovisioned resource从而兼容桶名不遵循amazon-sagemaker-{account_id}-{region}-{project_id}模板的 BYOR-bucket 项目算子用S3Hook.read_key读取该 JSON把每个键值对以NOTEBOOK_OUTPUT.{key}为 key 推送到 XCom并同时推送notebook_run_id。与 Sensor / Trigger 的配合SensorSageMakerUnifiedStudioNotebookSensor 轮询 DataZoneGetNotebookRunAPI直到运行进入终态SUCCEEDED成功QUEUED/STARTING/RUNNING/STOPPING视为进行中其他状态抛错。它接受notebook_run_id来自算子的输出以及notebook_identifier用于完成后读取 S3 输出并推送 XCom。示例run_sensor SageMakerUnifiedStudioNotebookSensor( task_idnotebook-sensor-task, aws_conn_idDATAZONE_CONN_ID, domain_identifierdomain_id, owning_project_identifierproject_id, notebook_identifiernotebook_id, notebook_run_idrun_notebook.output[notebook_run_id], )注意这里把算子的返回值run_notebook.output[notebook_run_id]直接作为 Sensor 的输入体现了 Airflow 3 中算子输出与下游任务的数据流衔接。TriggerSageMakerUnifiedStudioNotebookTrigger 使用自定义 boto waiternotebook_run_complete定义于 waiters 配置中异步监控运行实现deferrableTrue时的可延后执行。另外需要注意版本前提DataZone NotebookRun API 要求botocore 1.43.1Hook 在每次调用前都会校验start_notebook_run/get_notebook_run方法是否存在缺失时给出明确的升级提示。如何选择两个算子的适用场景场景推荐算子执行项目内的 Jupyter notebook、querybook、Visual ETL 作业按文件相对路径定位工件SageMakerNotebookOperator执行已注册的 SageMaker Unified Studio Notebook按Notebook ID定位并需要输出变量跨任务传递SageMakerUnifiedStudioNotebookOperator选择时的其他考量若运行在 SageMaker Unified Studio 内置的 MWAA 环境中SageMakerNotebookOperator可直接利用环境注入的AIRFLOW__WORKFLOWS__*变量零配置运行在外部 Airflow 中则建议显式传入domain_id/project_id/domain_region要求sagemaker-studio1.0.25。若你的 Airflow 实例通过标准 AWS 连接aws_conn_id管理凭证且需要细粒度的角色代入如 DataZone 环境角色SageMakerUnifiedStudioNotebookOperator与AwsBaseOperator体系的集成更自然。两个算子都支持deferrable可延后模式与wait_for_completion控制长耗时 Notebook 建议启用 deferrable 以释放 worker 资源。深入阅读仓库中的相关实现与测试关联文档providers/amazon/docs/operators/sagemakerunifiedstudio.rst算子源码operators/sagemaker_unified_studio.py、operators/sagemaker_unified_studio_notebook.pyHook 源码hooks/sagemaker_unified_studio.py、hooks/sagemaker_unified_studio_notebook.py配套组件sensors/sagemaker_unified_studio_notebook.py、triggers/sagemaker_unified_studio_notebook.py、triggers/sagemaker_unified_studio.py、utils/sagemaker_unified_studio.py可运行的系统测试 DAGexample_sagemaker_unified_studio.py、example_sagemaker_unified_studio_notebook.py其中包含本文所有代码示例的完整上下文单元测试tests/unit/amazon/aws/operators/、tests/unit/amazon/aws/hooks/、tests/unit/amazon/aws/sensors/、tests/unit/amazon/aws/triggers/ 下的test_sagemaker_unified_studio*文件通过本文的算子与参数说明加上仓库中可直接参考的系统测试 DAG你可以快速把 SageMaker Unified Studio 的 Notebook、Querybook、Visual ETL 与 Unified Studio Notebook 执行纳入 Airflow 的调度与监控体系并借助 XCom 实现 Notebook 之间的数据传递。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考