Amazon MWAA Serverless 新增 PythonOperator 和 BashOperator 支持
DataHot 速览
Amazon Managed Workflows for Apache Airflow (MWAA) Serverless 现已支持 PythonOperator 和 BashOperator,用户可直接在无服务器运行时中运行自定义 Python 函数和 shell 脚本,无需再借助 Lambda 或 ECS 等额外计算服务。该功能可降低编排管线的复杂性、成本和延迟,适用于 ETL 和数据质量检查等场景。文章还演示了将 CSV 转换为 JSON 的示例流程。
为什么值得关注:数据工程师可据此简化 Airflow 无服务器编排架构,减少对外部计算服务的依赖,值得关注。
本文目录 19 节
译文
AI 逐段翻译如果你在Apache Airflow 工作流上使用Amazon MWAA Serverless,现在可以使用 PythonOperator 和 BashOperator 在无服务器运行时直接运行自定义代码。此前,Amazon Managed Workflows for Apache Airflow (Amazon MWAA) Serverless 仅支持通过操作符编排 AWS 服务,用于调度任务、管理依赖和处理重试。它不支持原生运行你自己的 Python 函数或 shell 脚本。如果你需要自定义 Python 逻辑或 shell 命令,则必须将代码包装在 AWS Lambda 函数中、启动 Amazon Elastic Container Service (Amazon ECS) 任务,或使用其他 AWS 计算服务。这些替代方案为你的编排管道增加了复杂性、成本和延迟。
通过此发布,你可以直接在无服务器任务运行时内运行自定义 Python 函数和 shell 脚本,无需额外的基础设施。这意味着你现在可以使用许多数据工程团队在 ETL 管道和数据质量检查中依赖的PythonOperator和BashOperator——无需预置额外计算资源。
在这篇文章中,我们将介绍此功能的工作原理,并演示一个实际示例:使用 PythonOperator 构建一个将 CSV 文件转换为 JSON 格式的无服务器管道,并使用 BashOperator 验证输出。最后,你将了解如何:
- 将 Python 模块及其依赖打包,并将其作为代码包上传到 Amazon Simple Storage Service (Amazon S3) 存储桶
- 使用dag-factory兼容的 YAML 定义多任务工作流
- 使用 AWS Command Line Interface (AWS CLI) 创建并运行工作流
- 验证你的管道是否产生了预期输出
工作原理
使用 MWAA Serverless,你可以打包自定义代码,将其上传到 Amazon S3 存储桶,并在创建工作流时引用它。服务在工作流创建时对你的代码进行快照,并在同一工作流版本的所有后续运行中使用该快照。
代码包
代码包是包含你自定义逻辑的包。你将 Python 模块或 shell 脚本打包并上传到 Amazon S3 存储桶。代码包可以是:
- 单个 .py 文件或 .sh bash 脚本(上传到 Amazon S3 存储桶)
- 包含多个 shell 脚本、Python 模块和依赖的 ZIP 归档(最大 250 MB)
执行模型
当你创建或更新工作流时,MWAA Serverless 从提供的 Amazon S3 存储桶中对你的代码包进行快照,并将其存储在服务端。在任务执行时,服务使用此快照——而不是当前存在于你 Amazon S3 存储桶中的对象——在隔离的运行时环境中运行你的代码。
Python 和 Bash 任务无法访问互联网。它们只能访问 Amazon S3、Amazon Elastic Container Registry (Amazon ECR) 和 Amazon CloudWatch,这些是运行时运行所需的服务。要获得互联网访问,请使用 Amazon VPC 配置工作流,以便它可以通过提供的 VPC。
支持的运算符
下表描述了 MWAA Serverless 中现在可用的两个运算符。
| 运算符 | 描述 |
| PythonOperator | 从你的代码包中执行 Python 可调用对象(函数) |
| BashOperator | 运行 shell 命令或脚本 |
安全性
AWS Key Management Service (AWS KMS) 在静态时加密你的代码包。IAM 策略控制谁可以创建、更新和触发工作流。执行角色限定你的代码在运行时可以访问哪些 AWS 资源。
前提条件
开始之前,请验证你的 AWS 账户中配置了以下资源和工具:
- 一个有权访问 Amazon MWAA Serverless 的 AWS 账户
- 已安装并配置 AWS CLI v2(最新版本)。要安装或升级,请参阅安装或升级到最新版本的 AWS CLI。
- 一个用于存储 DAG 定义和代码包的 Amazon S3 存储桶
- MWAA Serverless 可以担任的 IAM 角色(请参阅下面的执行角色设置)
演练:构建无服务器 CSV 到 JSON 管道
在此演练中,我们构建一个将 CSV 文件转换为 JSON 格式的管道——这是下游 API 和分析系统消费 JSON 的常见数据转换。该管道使用 PythonOperator 进行转换逻辑,并使用 BashOperator 验证输出。管道执行以下操作:
- 从 Amazon S3 存储桶读取 CSV 文件
- 将其转换为 JSON 格式,并进行列类型推断
- 将 JSON 文件写回 Amazon S3 存储桶
- 验证源文件和输出之间的记录数是否匹配
步骤 1:创建执行角色
创建一个 IAM 角色,你的工作流在运行时担任。信任策略必须允许 airflow-serverless.amazonaws.com 服务担任该角色:
cat > trust-policy.json << 'EOF'
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Principal": {
"Service": "airflow-serverless.amazonaws.com"
},
"Action": "sts:AssumeRole"
}
]
}
EOF创建角色并附加内联策略,授予对你的 S3 存储桶的最低权限访问权限:
aws iam create-role \
--role-name MWAAServerlessExecutionRole \
--assume-role-policy-document file://trust-policy.json
aws iam put-role-policy \
--role-name MWAAServerlessExecutionRole \
--policy-name MWAAServerlessAccessPolicy \
--policy-document '{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:ListBucket"
],
"Resource": [
"arn:aws:s3:::amzn-s3-demo-mwaa-data",
"arn:aws:s3:::amzn-s3-demo-mwaa-data/*"
]
},
{
"Effect": "Allow",
"Action": [
"logs:CreateLogGroup",
"logs:CreateLogStream",
"logs:PutLogEvents",
"logs:DescribeLogStreams",
"logs:GetLogEvents"
],
"Resource": "arn:aws:logs:*:*:log-group:/aws/mwaa-serverless/*"
}
]
}'步骤 2:编写 Python 模块
创建一个名为csv_to_json.py的文件,包含转换逻辑:
# csv_to_json.py
import csv
import json
import boto3
import io
def convert(**kwargs):
"""Read a CSV from S3 and write it back as JSON lines."""
bucket = "amzn-s3-demo-mwaa-data"
source_key = "raw/sales_data.csv"
output_key = "processed/sales_data.json"
s3 = boto3.client("s3")
# Read source file
response = s3.get_object(Bucket=bucket, Key=source_key)
content = response["Body"].read().decode("utf-8")
# Parse CSV
reader = csv.DictReader(io.StringIO(content))
rows = list(reader)
# Type inference - convert numeric fields
for row in rows:
for key, value in row.items():
try:
row[key] = float(value)
except (ValueError, TypeError):
pass
# Write as JSON lines
output = "\n".join(json.dumps(row) for row in rows) + "\n"
s3.put_object(Bucket=bucket, Key=output_key, Body=output.encode("utf-8"))
print(f"Converted {len(rows)} rows to JSON lines")
print(f"Output: s3://amzn-s3-demo-mwaa-data/{output_key}")
return {"rows": len(rows), "output_key": output_key}此函数使用boto3(随 MWAA Serverless 执行环境预装)以及 Python 的内置 csv 和 json 模块。转换读取 CSV,推断数字类型,并将 JSON 行文件写回 S3 存储桶。
步骤 3:编写验证脚本
创建一个名为verify_output.sh的文件。此脚本通过比较源 CSV 中的记录数与输出 JSON 文件来验证管道输出。如果计数不匹配,任务将以非零退出码失败,导致工作流运行失败。
#!/bin/bash
echo "=== Data Validation ==="
# Count source records (skip CSV header)
SOURCE_COUNT=$(python3 -m awscli s3 cp s3://amzn-s3-demo-mwaa-data/raw/sales_data.csv - | tail -n +2 | wc -l)
echo "Source CSV records: $SOURCE_COUNT"
# Count output records
OUTPUT_COUNT=$(python3 -m awscli s3 cp s3://amzn-s3-demo-mwaa-data/processed/sales_data.json - | wc -l)
echo "Output JSON records: $OUTPUT_COUNT"
# Validate counts match
if [ "$SOURCE_COUNT" -ne "$OUTPUT_COUNT" ]; then
echo "FAILED: Record count mismatch (source=$SOURCE_COUNT, output=$OUTPUT_COUNT)"
exit 1
fi
echo "PASSED: Record counts match ($OUTPUT_COUNT records)"
echo "Timestamp: $(date -u +%Y-%m-%dT%H:%M:%SZ)"此脚本运行 AWS CLI,该 CLI 作为依赖项捆绑在代码包中。s3 cp 将文件内容流式传输到stdout,而不写入磁盘,允许标准 shell 工具(如wc -l和tail)处理它。执行角色凭据在执行环境中自动可用,因此 CLI 无需额外配置即可访问 S3。
步骤 4:打包并上传代码到 Amazon S3
由于验证脚本使用 AWS CLI,请将其作为依赖项与 Python 模块和 shell 脚本一起捆绑在 ZIP 归档中:
BUCKET="amzn-s3-demo-mwaa-data"
REGION="us-east-1"
# Install awscli into a package directory
pip install awscli \
--target my_package/ \
--platform manylinux2014_x86_64 \
--python-version 3.12 \
--only-binary=:all:
# Add your module
cp csv_to_json.py my_package/
cp verify_output.sh my_package/
# Create the ZIP archive
cd my_package && zip -r ../code_bundle.zip . && cd ..
# Upload to S3
aws s3 cp code_bundle.zip s3://$BUCKET/code/code_bundle.zip --region $REGION上传一个示例 CSV 文件用于测试:
cat > sales_data.csv << 'EOF'
date,region,product,units,revenue
2026-07-01,us-east,widget-a,150,4500.00
2026-07-01,eu-west,widget-b,89,2670.00
2026-07-02,us-east,widget-a,203,6090.00
2026-07-02,ap-south,widget-c,67,1340.00
2026-07-03,us-east,widget-b,178,5340.00
EOF
aws s3 cp sales_data.csv s3://$BUCKET/raw/sales_data.csv --region $REGION步骤 5:定义 DAG (YAML)
MWAA Serverless 使用声明式 YAML 格式定义 DAG。创建一个名为conversion_dag.yaml的文件:
csv_to_json_pipeline:
start_date: "2026-01-01"
schedule: null
tasks:
convert_to_json:
operator: airflow.operators.python.PythonOperator
python_callable: csv_to_json.convert
verify_output:
operator: airflow.operators.bash.BashOperator
bash_command: "verify_output.sh"
dependencies:
- convert_to_json此 DAG 定义了两个任务:
convert_to_json– 运行 Python 模块中的 convert 函数,将 CSV 转换为 JSON 行。verify_output– 运行一个 shell 脚本,通过比较源记录数和输出记录数来验证管道输出,如果不匹配则任务失败。
将 DAG 定义上传到 S3。注意:你也可以直接运行内联 Bash 命令,无需 shell 脚本。
aws s3 cp conversion_dag.yaml s3://$BUCKET/dags/conversion_dag.yaml --region $REGION步骤 6:创建工作流
创建 MWAA Serverless 工作流,引用 DAG 定义和代码包:
ROLE_ARN="arn:aws:iam::<your-account-id>:role/MWAAServerlessExecutionRole"
aws mwaa-serverless create-workflow \
--name csv-to-json-workflow \
--definition-s3-location Bucket="$BUCKET",ObjectKey="dags/conversion_dag.yaml" \
--code '{"S3Location": {"Bucket":"'"$BUCKET"'","ObjectKey":"code/code_bundle.zip"}}' \
--role-arn $ROLE_ARN \
--region $REGION响应中包含一个 WorkflowArn,用于触发运行:
{
"WorkflowArn": "arn:aws:airflow-serverless:us-east-1:123456789012:workflow/csv-to-json-workflow-abc123",
"CreatedAt": "2026-07-15T10:30:00.000000+00:00",
"WorkflowVersion": "a1b2c3d4e5f6"
}步骤 7:运行工作流
触发工作流运行:
WORKFLOW_ARN="arn:aws:airflow-serverless:us-east-1:123456789012:workflow/csv-to-json-workflow-abc123"
aws mwaa-serverless start-workflow-run \
--workflow-arn $WORKFLOW_ARN \
--region $REGION响应确认运行已开始:
{
"RunId": "6OZV9ABF9enHKXk",
"Status": "STARTING"
}步骤 8:监控执行
检查运行状态:
RUN_ID="6OZV9ABF9enHKXk"
aws mwaa-serverless get-workflow-run \
--workflow-arn $WORKFLOW_ARN \
--run-id $RUN_ID \
--region $REGION成功运行返回:
{
"RunDetail": {
"Duration": 45,
"RunState": "SUCCESS",
"TaskInstances": ["ex_abc123_convert_to_json_1", "ex_abc123_verify_output_1"]
},
"RunId": "6OZV9ABF9enHKXk",
"RunType": "ON_DEMAND",
"WorkflowArn": "arn:aws:airflow-serverless:us-east-1:123456789012:workflow/csv-to-json-workflow-abc123",
"WorkflowVersion": "a1b2c3d4e5f6"
}步骤 9:验证输出
确认 JSON 文件已写入 S3 存储桶:
# List the output file
aws s3 ls s3://$BUCKET/processed/sales_data.json --region $REGION你应该能看到 JSON 文件:
2026-07-15 10:32:45 1847 sales_data.json你还可以在 Amazon CloudWatch Logs 中验证任务级输出。打开工作流的日志组,找到convert_to_json任务日志流:
Converted 5 rows to JSON lines
Output: s3://amzn-s3-demo-mwaa-data/processed/sales_data.json注意事项和限制
在使用这些运算符规划 MWAA Serverless 上的工作负载时,请牢记以下注意事项:
- 代码包大小 – ZIP 归档必须小于 250 MB。
- 网络访问 – Python 和 Bash 任务没有互联网访问权限。它们只能访问运行所需的一组有限的 AWS 服务(Amazon S3、Amazon ECR 和 Amazon CloudWatch),但不能调用其他 AWS 服务或外部端点。如果你的工作流需要调用外部 API,请在工作流调用之前预处理数据并将其存储在 Amazon S3 存储桶中。
- 运行时依赖 – 预装了 boto3 和 Python 标准库。对于其他包(如 pandas 或 requests),请按照Amazon MWAA Serverless 打包指南将它们打包到你的 ZIP 归档中。
- 执行超时 – 任务受工作流配置的超时限制。
- Python 版本 – 查看Amazon MWAA Serverless 文档了解当前支持的 Python 运行时版本。
- DAG 格式 – MWAA Serverless 使用基于 YAML 的 DAG 定义,而不是传统的 Python DAG 文件。如果从 MWAA Provisioned 迁移,你需要将 DAG 转换为 YAML 格式。
- 不支持的运算符 – 一些 Airflow 社区运算符和自定义插件在 Serverless 运行时中不可用。请参阅文档获取完整的兼容性列表。
清理
为避免持续产生费用,请删除本演练中创建的资源。以下命令删除工作流、S3 对象和 IAM 角色:
注意:$WORKFLOW_ARN 在步骤 7 中定义。
# Delete the workflow
aws mwaa-serverless delete-workflow \
--workflow-arn $WORKFLOW_ARN \
--region $REGION注意:$BUCKET 在步骤 4 中导出。如果合适,也删除存储桶。
# Remove S3 objects
aws s3 rm s3://$BUCKET/code/code_bundle.zip
aws s3 rm s3://$BUCKET/dags/conversion_dag.yaml
aws s3 rm s3://$BUCKET/raw/sales_data.csv
aws s3 rm s3://$BUCKET/processed/sales_data.json# Delete the IAM role
aws iam delete-role-policy \
--role-name MWAAServerlessExecutionRole \
--policy-name MWAAServerlessAccessPolicy
aws iam delete-role --role-name MWAAServerlessExecutionRole结论
凭借对 PythonOperator 和 BashOperator 的原生支持,你现在可以在 MWAA Serverless 中直接运行许多数据工程团队日常依赖的自定义代码执行模式。在无服务器运行时中运行数据转换、格式转换、验证和 shell 脚本 – 无需预置额外计算或管理容器。
如果你正在 MWAA Provisioned 或自管理基础设施上运行 Airflow 工作负载,你现有的 PythonOperator 和 BashOperator 逻辑只需进行少量更改。将 Python DAG 文件转换为 YAML 格式,将代码打包为包,即可在 MWAA Serverless 上运行。
要开始使用,请访问Amazon MWAA Serverless 文档,并尝试本文前面的演练,使用你自己的数据。有关定价详情,请访问Amazon MWAA 定价页面。我们期待您的反馈。
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏