Amazon MWAA ServerlessのPythonOperatorで複数のSSMパラメーターを1タスクで取得してみた

Amazon MWAA ServerlessのPythonOperatorで複数のSSMパラメーターを1タスクで取得してみた

MWAA ServerlessのPythonOperatorでSSMの複数パラメーターを一括取得し、辞書XComのブラケットアクセスで下流タスクへ渡す構成を確認しました。タスク分割は不要になりますが、VPCエンドポイントの固定費を含めるとLambda回避策やDynamoDB構成のほうが安いケースもあります。
2026.09.22

こんにちは。サービス開発部の武田です。

Amazon MWAA Serverlessのワークフロー定義はYAMLで書きます。バケット名や出力先接頭辞のような設定値をAWS Systems Manager Parameter Storeに置いて実行時に取得したい、という場面は多いはずです。しかしこれまでのMWAA Serverlessでは、複数のパラメーターをまとめて取得する簡単な方法がありませんでした。Lambda経由で取得してもLambdaInvokeFunctionOperatorがレスポンスを文字列としてXComに格納するため、テンプレートからキーで取り出せません。パラメーターごとにLambdaタスクを分けるなどの回避策が必要でした。回避策の詳細は次の記事にまとめています。

https://dev.classmethod.jp/articles/amazon-mwaa-serverless-lambda-multi-value-workarounds/

2026年8月に、MWAA ServerlessでPythonOperatorとBashOperatorが使えるようになりました。

https://aws.amazon.com/about-aws/whats-new/2026/08/mwaa-serverless-pythonoperator-bashoperator/

自前のPython関数でboto3を呼べるなら、get_parametersで複数パラメーターを一括取得して、辞書としてまとめてXComに載せられるはずです。動かして確認しました。

PythonOperatorの基本的な使い方は次の記事を参照してください。

https://dev.classmethod.jp/articles/update-amazon-mwaa-serverless-support-pythonoperator-bashoperator/

辞書を返すと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サービスの範囲は、次の記事で確認しています。

https://dev.classmethod.jp/articles/amazon-mwaa-serverless-vpc-network-reachability/

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で失敗を診断できる設計にしておく

どなたかの参考になりましたら幸いです。

参考リンク

この記事をシェアする

AWSのお困り事はクラスメソッドへ

関連記事