Amazon MWAA ServerlessのPythonOperatorで複数のSSMパラメーターを1タスクで取得してみた
こんにちは。サービス開発部の武田です。
Amazon MWAA Serverlessのワークフロー定義はYAMLで書きます。バケット名や出力先接頭辞のような設定値をAWS Systems Manager Parameter Storeに置いて実行時に取得したい、という場面は多いはずです。しかしこれまでのMWAA Serverlessでは、複数のパラメーターをまとめて取得する簡単な方法がありませんでした。Lambda経由で取得してもLambdaInvokeFunctionOperatorがレスポンスを文字列としてXComに格納するため、テンプレートからキーで取り出せません。パラメーターごとにLambdaタスクを分けるなどの回避策が必要でした。回避策の詳細は次の記事にまとめています。
2026年8月に、MWAA ServerlessでPythonOperatorとBashOperatorが使えるようになりました。
自前のPython関数でboto3を呼べるなら、get_parametersで複数パラメーターを一括取得して、辞書としてまとめてXComに載せられるはずです。動かして確認しました。
PythonOperatorの基本的な使い方は次の記事を参照してください。
辞書を返すとXComも辞書型になる
ポイントは、python_callableが辞書を返すと、XComには辞書型のまま格納されることです。MWAA Serverlessのテンプレートはドットアクセス(.key)に制限がありますが、ブラケットアクセス(['key'])は動作します。辞書型のXComであれば、下流タスクのテンプレートから個別の値を取り出せます。
まず、この動きだけを切り出して確認しました。合わせて、タスクからSSMに到達できるかも見ておきます。コードは次のとおりです。
# probe.py
import boto3
from botocore.config import Config
CFG = Config(connect_timeout=3, read_timeout=5, retries={"max_attempts": 0})
def get_config():
try:
boto3.client("ssm", config=CFG).get_parameter(Name="/myapp/dummy")
ssm = "ok"
except Exception as e:
ssm = f"{e.__class__.__name__}"
result = {"ssm": ssm, "data_bucket": "my-bucket", "output_prefix": "data/"}
print(f"result = {result}")
return result
ワークフロー定義は次のとおりです。下流のBashOperatorで、3つのキーをブラケットアクセスで取り出します。
pybash_ssmprobe:
dag_id: pybash_ssmprobe
schedule: null
default_args:
owner: airflow
start_date: "2024-01-01"
tasks:
get_config:
operator: airflow.providers.standard.operators.python.PythonOperator
task_id: get_config
python_callable: probe.get_config
bash_show:
operator: airflow.providers.standard.operators.bash.BashOperator
task_id: bash_show
bash_command: "echo \"bucket={{ ti.xcom_pull(task_ids='get_config')['data_bucket'] }} prefix={{ ti.xcom_pull(task_ids='get_config')['output_prefix'] }} ssm={{ ti.xcom_pull(task_ids='get_config')['ssm'] }}\""
dependencies:
- get_config
VPCなし(NetworkConfiguration未指定)で実行すると、ワークフローはSUCCESSになり、bash_showのレンダリング結果は次のようになりました。
bucket=my-bucket prefix=data/ ssm=ConnectTimeoutError
2つのことが分かります。
- 辞書型XComへのブラケットアクセスは、PythonOperatorの戻り値でも動作する(
bucket=とprefix=が展開されている) - VPCなしのタスクからSSMには到達できない(
ConnectTimeoutError)
SSMに届くには自前VPCが必要です。VPCなしのタスクから到達できるAWSサービスの範囲は、次の記事で確認しています。
VPCとエンドポイントを用意する
NATゲートウェイなしのプライベートルーティング構成にしました。用意したVPCエンドポイントは次の3つです。
- S3ゲートウェイエンドポイント(定義とコードの取得に必要)
- SSMインタフェースエンドポイント(今回の呼び先)
- CloudWatch Logsインタフェースエンドポイント(タスクログを出すのに必要)
VPC本体は10.1.0.0/16に別AZのプライベートサブネット2つ、自己参照のインバウンドを許可したセキュリティグループという構成です。SSMのインタフェースエンドポイントは次のコマンドで作りました。
# SSM(インタフェースエンドポイント、プライベートDNS有効)
aws ec2 create-vpc-endpoint --vpc-id vpc-xxxxxxxx \
--service-name com.amazonaws.ap-northeast-1.ssm --vpc-endpoint-type Interface \
--subnet-ids subnet-private-a subnet-private-c \
--security-group-ids sg-xxxxxxxx --private-dns-enabled
create-vpcで作った直後のVPCは、DNSホスト名(enableDnsHostnames)が無効になっています。この状態で--private-dns-enabledを付けてエンドポイントを作成するとInvalidParameterで失敗します。先にmodify-vpc-attributeでDNS属性を有効化してください。
aws ec2 modify-vpc-attribute --vpc-id vpc-xxxxxxxx --enable-dns-hostnames
aws ec2 modify-vpc-attribute --vpc-id vpc-xxxxxxxx --enable-dns-support
実行ロールには、S3・CloudWatch Logsの権限に加えて、対象パラメーターへの読み取り権限を追加します。
{
"Sid": "SsmGetParameters",
"Effect": "Allow",
"Action": ["ssm:GetParameters"],
"Resource": "arn:aws:ssm:ap-northeast-1:123456789012:parameter/myapp/*"
}
テスト用のパラメーターは、String型で3つ作りました。
| パラメーター名 | 値 |
|---|---|
/myapp/data_bucket |
amzn-s3-demo-bucket |
/myapp/output_prefix |
definitions/ |
/myapp/table_name |
my-config-table |
複数パラメーターを1タスクで取得する
取得コードは次のとおりです。get_parametersで3つを一括取得し、パラメーター名の接頭辞を外した辞書にして返します。
# ssmcfg.py
import boto3
PREFIX = "/myapp/"
NAMES = ["data_bucket", "output_prefix", "table_name"]
def get_config():
ssm = boto3.client("ssm")
resp = ssm.get_parameters(Names=[PREFIX + n for n in NAMES])
result = {p["Name"].removeprefix(PREFIX): p["Value"] for p in resp["Parameters"]}
print(f"result = {result}")
return result
ワークフロー定義は3タスク構成です。list_files(S3ListOperator)がバケット名と接頭辞を、bash_showがテーブル名を、それぞれブラケットアクセスで受け取ります。
pybash_ssmcfg:
dag_id: pybash_ssmcfg
schedule: null
default_args:
owner: airflow
start_date: "2024-01-01"
tasks:
get_config:
operator: airflow.providers.standard.operators.python.PythonOperator
task_id: get_config
python_callable: ssmcfg.get_config
list_files:
operator: airflow.providers.amazon.aws.operators.s3.S3ListOperator
task_id: list_files
bucket: "{{ ti.xcom_pull(task_ids='get_config')['data_bucket'] }}"
prefix: "{{ ti.xcom_pull(task_ids='get_config')['output_prefix'] }}"
dependencies:
- get_config
bash_show:
operator: airflow.providers.standard.operators.bash.BashOperator
task_id: bash_show
bash_command: "echo \"table={{ ti.xcom_pull(task_ids='get_config')['table_name'] }} first_file={{ ti.xcom_pull(task_ids='list_files')[0] }}\""
dependencies:
- list_files
作成時に--network-configurationで先ほどのVPCを渡します。
aws mwaa-serverless create-workflow \
--name ssm-config \
--definition-s3-location '{"Bucket":"amzn-s3-demo-bucket","ObjectKey":"definitions/ssm-config.yaml"}' \
--code '{"S3Location":{"Bucket":"amzn-s3-demo-bucket","ObjectKey":"code/ssm-config.zip"}}' \
--role-arn arn:aws:iam::123456789012:role/mwaa-serverless-exec \
--network-configuration '{"SecurityGroupIds":["sg-xxxxxxxx"],"SubnetIds":["subnet-private-a","subnet-private-c"]}'
実行すると、3タスクすべてSUCCESSで完走しました。get_configのタスクログに一括取得の結果が出ています。
result = {'data_bucket': 'amzn-s3-demo-bucket', 'output_prefix': 'definitions/', 'table_name': 'my-config-table'}
list_filesのXComには、取得したバケット名と接頭辞でS3をリストした結果が入っていました。bash_showの出力も期待どおりです。
table=my-config-table first_file=definitions/ssm-config.yaml
パラメーターごとのLambdaタスク分割も、固定幅パディングも、Step Functions経由も不要になりました。パラメーターが増えても、タスクは1つのままです。
VPCエンドポイントの時間課金
回避策で一番問題だったのは、MWAA Serverlessの課金がタスクインスタンスごとに最低1分であることでした。パラメーターごとにタスクを分けると、パラメーター数に比例して課金時間が増えます。1タスクにまとめられたので、パラメーターが増えても課金時間は増えません。
代わりに、VPCエンドポイントの時間課金がかかります。インタフェースエンドポイントは東京だと0.014 USD/時/AZです。SSMとCloudWatch Logsの2つを2AZに置くと約41 USD/月になります。ワークフロー専用にVPCを新設するならこのコストが乗りますが、既存のVPCにエンドポイントがすでにあるなら追加コストはほぼありません。
ストアをDynamoDBに変える選択肢
SSM構成で残るコストはインタフェースエンドポイントの時間課金です。設定値のストアをDynamoDBに変えると、SSMのぶんを減らせます。DynamoDBはS3と同じく ゲートウェイエンドポイント に対応しており、ゲートウェイエンドポイントには課金がありません。
DynamoDBに変えてもVPC自体は必要です。前述の到達性の記事で確認したとおり、VPCなしのタスクからはDynamoDBにも到達できません。変わるのはエンドポイントの種類と課金です。
設定値は1アイテムに複数属性で持たせます。
| 属性 | 値 |
|---|---|
config_id(パーティションキー) |
app |
data_bucket |
amzn-s3-demo-bucket |
output_prefix |
definitions/ |
table_name |
my-config-table |
取得コードではget_itemで1アイテムを読み、パーティションキー以外の属性を辞書にして返します。
# ddbcfg.py
import boto3
def get_config():
try:
table = boto3.resource("dynamodb").Table("my-config-table")
item = table.get_item(Key={"config_id": "app"})["Item"]
result = {k: v for k, v in item.items() if k != "config_id"}
except Exception as e:
result = {"error": f"{e.__class__.__name__}: {e}"}
print(f"result = {result}")
return result
ポイントは2点です。
boto3.resourceを使う。低レベルのclientだと属性値が{"S": "..."}の型付き形式になり、テンプレートからの取り出しが['data_bucket']['S']と2段になる。resourceならプレーンな文字列の辞書になる- 例外を握って辞書に包んで返す。理由は後述
実行ロールにはdynamodb:GetItemを対象テーブルに絞って追加します。
VPCエンドポイントをS3ゲートウェイとDynamoDBゲートウェイの2つだけにした構成、つまりインタフェースエンドポイントが1つもない構成で実行しました。結果は3タスクすべてSUCCESSで、SSM版と同じようにブラケットアクセスで値を取り出せました。
ただし、CloudWatch Logsのエンドポイントがないため、 タスクログが残りません。ロググループ自体は作られます(ロググループ名を指定しなければ/aws/mwaa-serverless/<ワークフロー名>/が使われるとドキュメントに記載があります)が、ログストリームは空のままでした。実行は成功するので、ログがないことに気付きにくいです。例外を辞書に包んでXComで返しているのはこのためです。ログが出ない構成でも、get-task-instanceのXComから失敗の内容を確認できます。
回避策だったパラメーターごとのLambdaタスクを含めて、4パターンの月額を比べます。前提はパラメーター3個で、下流の共通タスクは除きます。タスク課金は0.08 USD/時(1分の最低課金で0.0013 USD/タスク)、インタフェースエンドポイントは0.014 USD/時/AZの2AZで730時間として計算しています。Lambda自体の料金、SSMの標準パラメーター、DynamoDBの読み取りは無視できる額なので入れていません。
| 構成 | インタフェースエンドポイント | タスクログ | 月額固定 | 1実行のタスク課金 | 月額(1日1回) | 月額(1時間に1回) |
|---|---|---|---|---|---|---|
| Lambda×3タスク(VPCなし) | なし | あり | 0 USD | 0.004 USD | 0.12 USD | 2.88 USD |
| SSM + Logs | 2つ | あり | 40.88 USD | 0.0013 USD | 40.92 USD | 41.84 USD |
| DynamoDB + Logs | 1つ | あり | 20.44 USD | 0.0013 USD | 20.48 USD | 21.40 USD |
| DynamoDBのみ | なし | なし | 0 USD | 0.0013 USD | 0.04 USD | 0.96 USD |
1実行あたりのタスク課金はLambda方式が3倍ですが、月額ではインタフェースエンドポイントの固定費のほうが大きいです。SSM + Logs構成がLambda方式より安くなるのは、パラメーター3個なら月15,000回(3分に1回)を超えてから、10個でも月3,400回(13分に1回)前後からです。既存VPCのエンドポイントを使い回せないなら、コストだけで選ぶ場合はLambda方式かDynamoDBのみの構成になります。
なお、SSMの「N個のパラメーター」に対してDynamoDBは「1アイテムのN属性」と、設定値の持ち方の粒度が変わります。パラメーター単位のアクセス制御や履歴管理といったParameter Storeの機能は使えなくなるので、単純なコスト比較だけで置き換えを決めないでください。
まとめ
- PythonOperatorの
python_callableが返す辞書は、辞書型のままXComに格納され、下流テンプレートのブラケットアクセスで取り出せる - これにより、複数のSSMパラメーターの一括取得が1タスクで書けるようになった。パラメーター数に比例したタスク課金はなくなったが、インタフェースエンドポイントの固定費のほうが大きく、パラメーター3個・1時間に1回の実行ならLambda方式のほうが安い
- SSMへの到達には自前VPCが必要。SSMインタフェースエンドポイントに加えて、タスクログ用のCloudWatch Logsエンドポイントも忘れないこと
- ストアをDynamoDBに変えると、課金なしのゲートウェイエンドポイントだけで構成できる。ただしLogsエンドポイントを省くとタスクログが残らないため、XComで失敗を診断できる設計にしておく
どなたかの参考になりましたら幸いです。






