インフラ

AWS Athenaの実装:SDKでのクエリ実行とフェデレーテッドクエリ・Iceberg更新

AWS Athenaの実装:SDKでのクエリ実行とフェデレーテッドクエリ・Iceberg更新

AWS Athenaをコンソールで触るところまでは、たいてい詰まりません。詰まるのはその次です。日次のバッチから叩こうとしてポーリングの書き方で止まり、S3の外にあるRDSのマスタと結合したくなって手が止まる。この記事では2026年9月22日時点のAWS公式ユーザーガイドをもとに、Athenaをアプリケーションとデータ基盤へ組み込む実装と、運用で効く制約を整理します。仕組みやクエリの書き方はAthenaの仕組みと使い方をまとめた解説に、請求額の分解はスキャン量課金の実額を扱った記事に譲ります。

まとめ:組み込みの要はSDKの非同期実行とワークグループ単位の分離

先に結論を置きます。Athenaをプログラムから使うのは、同期的なデータベース接続を書くこととは別物です。start_query_execution はクエリIDを返すだけで結果は返しません。状態をポーリングし、完了を確認してから結果を取りに行く。この非同期の形を受け入れられるかが、組み込みの成否を分けます。

もう一つの軸がワークグループです。AWS公式ユーザーガイドはワークグループをIAMリソースとして扱い、ワークロードの分離とコスト制御に使うと説明しています。エンジンバージョンの固定も、スキャン量の上限も、結果の再利用も、すべてワークグループ単位で決まる設定です。アプリ用とアナリスト用の同居は後で必ず壊れます。

機能面の上限も先に知っておくべきでしょう。更新と削除が要るならIcebergテーブルですが、Athenaが作れるのはv2だけで、タイムスタンプはミリ秒精度に丸められます。S3の外を混ぜるならフェデレーテッドクエリですが、Lambdaの実行料金が別建てで乗ります。

boto3でクエリを投げて結果を受け取るまでの非同期実装とエラー処理

Athenaへの入口はコンソール、JDBC/ODBCドライバ、AWS SDKの3つで、公式ユーザーガイドのサービス概要が示すとおりどの経路でも起動しておくクラスターはありません。コンソールは検証の場であって定期実行の置き場ではなく、ドライバ経由はポーリングが隠れる代わりに版数が機能の可否を決めます。組み込むならSDKです。Pythonならboto3の start_query_execution でクエリを投入し、get_query_execution で状態を見て、get_query_results で取り出します。間の待ち方に実装差が出ます。

start_query_execution からポーリングまでのboto3実装の全体像

最小限の形を示します。クエリIDを受け取り、終端まで待ち、成功なら結果を読む流れです。

import time
import boto3

athena = boto3.client("athena", region_name="ap-northeast-1")

def run_query(sql, workgroup="app-batch", database="analytics"):
    started = athena.start_query_execution(
        QueryString=sql,
        QueryExecutionContext={"Database": database},
        WorkGroup=workgroup,
        ResultReuseConfiguration={
            "ResultReuseByAgeConfiguration": {
                "Enabled": True,
                "MaxAgeInMinutes": 60,
            }
        },
    )
    qid = started["QueryExecutionId"]

    wait = 0.5
    deadline = time.time() + 900
    while True:
        res = athena.get_query_execution(QueryExecutionId=qid)
        status = res["QueryExecution"]["Status"]
        state = status["State"]
        if state in ("SUCCEEDED", "FAILED", "CANCELLED"):
            break
        if time.time() > deadline:
            athena.stop_query_execution(QueryExecutionId=qid)
            raise TimeoutError(f"athena query timeout: {qid}")
        time.sleep(wait)
        wait = min(wait * 2, 5.0)

    if state != "SUCCEEDED":
        reason = status.get("StateChangeReason", "")
        raise RuntimeError(f"athena query {state}: {reason}")
    return qid

実行の状態は QUEUEDRUNNINGSUCCEEDEDFAILEDCANCELLED を取り、終端の3つ以外はまだ動いているという意味です。WorkGroup を省略した実装をよく見かけますが、省くと primary に落ちます。エンジンバージョンもスキャン量上限も primary の設定で動き、分離の設計が効かなくなります。

QUEUED と RUNNING を待つポーリング間隔とタイムアウトの決め方

上の例で待ち時間を0.5秒から5秒へ指数的に伸ばしているのには理由があります。Athenaのクォータはクエリの同時実行数だけでなく、APIの呼び出し回数にも掛かるからです。公式は秒あたり呼び出し数やバースト容量を超えると ThrottlingException が返ると明記しています。固定0.1秒でポーリングする実装は、本数が増えた時点でここに当たります。

QUEUED が長く続くときは、クエリが重いのではなく順番待ちです。公式のサービスクォータのページは、Active DMLとActive DDLのクォータに実行中とキュー待ちの両方が含まれると説明しています。DMLの同時実行クォータが25なら、合計26本目は TooManyRequestsException になります。DMLのタイムアウトは申請で最大240分まで伸ばせますが、伸びるのは1本あたりの実行時間であって、同時に走らせられる本数ではありません。

get_query_results のページングと型情報の扱いで詰まる箇所

get_query_results は1回で全件を返さず、NextToken を辿るページングが要ります。加えて先頭行がヘッダー行として返るため、そのまま処理すると列名がデータに混ざります。

def fetch_rows(qid):
    paginator = athena.get_paginator("get_query_results")
    first = True
    for page in paginator.paginate(QueryExecutionId=qid):
        rows = page["ResultSet"]["Rows"]
        if first:
            rows = rows[1:]  # 先頭はヘッダー行
            first = False
        for r in rows:
            yield [c.get("VarCharValue") for c in r["Data"]]

返される値の型は、すべて文字列です。ResultSetMetadata に列の型が入っているので、数値や日時に戻すなら自分でキャストします。NULLは VarCharValue キーそのものが存在しない形になるため、直接添字で書いた実装はNULLを含む列に当たった瞬間 KeyError で落ちます。上の例で get を使っているのはそのためです。件数が多いなら、結果ファイルのS3パスを get_query_execution から取り出して直接読むほうが速くなります。

S3の外にあるRDSやDynamoDBを横断するフェデレーテッドクエリの構成

ログはS3にあるが、顧客マスタはRDSにある。ETLを組まずに結合したいときに使うのがフェデレーテッドクエリです。データソースごとにLambda関数(コネクタ)を置き、Athenaがそれを呼び出して読みます。

Lambdaコネクタを介してRDSやDynamoDBへ接続する仕組み

公式のデータソース接続の手順では、コネクタと接続先の情報を保持するAWS Glueの接続を作り、そこで付けた名前をSQLの中でカタログ名として参照します。接続の作成方法は、コンソールからの操作と CreateDataCatalog APIの呼び出しのどちらでも選択可能です。VPC内のRDSにつなぐなら、Lambdaをそのサブネットに置く設定が別途要ります。テーブル定義そのものはGlue Data Catalogが持つため、クローラで作る運用ならGlueのETLジョブとDPU課金の設計と地続きの話になります。

SQLの側は、カタログ名をテーブル参照の先頭に付けるだけです。S3上のテーブルとキー設計と容量モードを扱ったDynamoDBのテーブルを同じクエリで結合できます。

SELECT
  l.request_id,
  l.status_code,
  u.plan_name
FROM awsdatacatalog.analytics.access_log AS l
JOIN "ddb-catalog".default.users AS u
  ON l.user_id = u.user_id
WHERE l.dt = '2026-09-21'
  AND l.status_code >= 500

フェデレーテッドクエリで別建てになるLambda料金と最小10MB課金

費用の数え方が通常のクエリと変わる点は、設計前に押さえておくべきでしょう。AWS公式の料金ページは、S3以外のデータソースへのクエリについて、横断で集計したスキャン量をTB単位で課金し、メガバイト単位に切り上げたうえでクエリごとに最小10MBを課金すると記載しています。加えてLambdaの実行料金が標準レートで別に掛かります。

ここから導かれる判断は明快です。数KBのマスタを引くためだけのフェデレーテッドクエリを1分おきに何百回も投げる構成は採りません。実データ量が10MBに満たなくても10MB分が課金され、そこにLambdaの起動が毎回乗るためです。小さなマスタなら日次でS3へ書き出すほうが安く速い。これが効くのは、アドホックな調査と、移し替えるには量が多すぎるデータの探索です。

Icebergテーブルでの更新・削除とタイムトラベルが使える条件と制約

S3上のテキストやParquetを普通のテーブルとして定義すると、行単位の更新も削除もできません。個人情報の削除要求に応えるため、パーティションごと作り直す運用を組んでいる現場は少なくないはずです。ここを解くのがIcebergテーブルになります。

Athenaが扱えるのはIceberg v2テーブルのみという前提と対応版1.4.2

公式のIcebergに関するページには前提がいくつも並んでいます。対応するApache Icebergのバージョンは1.4.2。作成・操作できるのはv2テーブルだけで、v1は対象外です。カタログはAWS Glueに限られ、ロックもGlueの楽観的ロックだけが対象。公式は、これ以外のロック実装で変更するとデータ損失とトランザクション破壊の可能性があると明記しています。

engine version 3 が読み書きできる形式はParquet、ORC、Avroの3つ。作成はCTASでもDDLでも書けます。

CREATE TABLE analytics.orders_iceberg (
  order_id   string,
  user_id    string,
  amount     double,
  ordered_at timestamp
)
PARTITIONED BY (day(ordered_at))
LOCATION 's3://example-lake/orders_iceberg/'
TBLPROPERTIES (
  'table_type' = 'ICEBERG',
  'format' = 'parquet'
);

MERGE INTO analytics.orders_iceberg AS t
USING analytics.orders_staging AS s
  ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET amount = s.amount
WHEN NOT MATCHED THEN INSERT (order_id, user_id, amount, ordered_at)
  VALUES (s.order_id, s.user_id, s.amount, s.ordered_at);

DELETE FROM analytics.orders_iceberg
WHERE user_id = 'u-000123';

パーティションに実装で引っかかる制約が一つ。ネストしたフィールドでのパーティションは対応しておらず、指定すると NOT_SUPPORTED: Partitioning by nested field is unsupported が返ります。JSONを構造体のまま持ち、その内側の日付で切る設計はここで作り直しです。

FOR TIMESTAMP AS OF への構文変更とミリ秒精度という2つの落とし穴

タイムトラベルは過去のある時点のテーブル内容をそのまま読む機能で、誤った更新を流した直後に変更前と突き合わせる用途で効きます。ただし構文が途中で変わりました。engine version 3 の変更点をまとめた公式ページは、以前の FOR SYSTEM_TIME AS OFFOR SYSTEM_VERSION AS OF が廃止され、FOR TIMESTAMP AS OFFOR VERSION AS OF になったと記載しています。古い記事を見て書くと mismatched input 'SYSTEM_TIME' で落ちます。

SELECT * FROM analytics.orders_iceberg
FOR TIMESTAMP AS OF (current_timestamp - interval '1' day);

SELECT * FROM analytics.orders_iceberg
FOR VERSION AS OF 949530903748831860;

もう一つが精度です。Icebergの仕様はtimestampにマイクロ秒精度を認めていますが、Athenaは読み書きともミリ秒精度しか扱わず、手動コンパクションで書き換えた時刻列も丸められます。マイクロ秒のデータをそのままCTASで入れると Incorrect timestamp precision for timestamp(6) で失敗するため、公式は一度 timestamp(6) へキャストして作り、参照側のビューで戻す回避策を示しています。ミリ秒未満の差で並び順が決まる計測データをAthenaで持つ設計は、ここで見直すべきでしょう。

MERGE と OPTIMIZE をLake Formationで権限管理できない制約

Icebergテーブルの運用にはメンテナンスが付いてきます。小さなファイルが増えるので OPTIMIZE でまとめ、不要なスナップショットを VACUUM で消す。この OPTIMIZE にはengine version 3特有の制約があり、WHERE 句に書けるのはパーティション列だけです。それ以外で絞ると Unexpected FilterNode found in plan が返ります。

権限設計にも穴があります。公式は、IcebergやHudi、Delta Lakeについて、Lake Formationで読み取り権限は管理できる一方、VACUUMMERGEUPDATEOPTIMIZE の権限は管理できないと明記しています。登録済みテーブルへのDDL操作自体も対象外です。Lake Formationによる権限管理で列・行レベルの制御を敷いた基盤に更新可能なIcebergテーブルを後から足すと、ここがぶつかります。

ワークグループでエンジンバージョンとコストと権限を分ける運用設計

ここまでの制約の多くが、ワークグループ単位の設定に紐づいています。アプリ用とアナリスト用を分けるのは綺麗事ではなく、機能の可否がワークグループで決まるという実装上の理由があるためです。クエリ結果の保存先もその一つで、書き出し先はS3のバケット作成と権限設計の手順で作った構成の上に置きます。ワークグループ側で保存先を強制するとクエリ単位の指定は無視され、逆にクエリ側で上書きすると結果の再利用が無効になります。

ワークグループ単位のエンジンバージョン固定と自動アップグレードの挙動

公式のエンジンバージョンのページは、バージョンがワークグループ単位で設定され、既定は自動アップグレードだと説明しています。自動なら、Athenaは非互換を検出しない限り勝手に上げます。明示したワークグループは指定のまま固定されますが、そのバージョンの提供が終わるときには上げられるため、永久の固定はできません。

本番バッチを載せるワークグループはバージョンを固定し、検証用は新しいバージョンで別に作って先に流す。公式が推奨する移行の形です。CLIなら次のように作れます。

aws athena create-work-group \
  --name app-batch \
  --configuration '{
    "ResultConfiguration": {"OutputLocation": "s3://example-athena-results/app-batch/"},
    "EnforceWorkGroupConfiguration": true,
    "EngineVersion": {"SelectedEngineVersion": "Athena engine version 3"},
    "BytesScannedCutoffPerQuery": 10737418240,
    "PublishCloudWatchMetricsEnabled": true
  }'

BytesScannedCutoffPerQuery は1クエリあたりのスキャン量の上限で、超えたクエリは実行中に自動キャンセルされます。絞り込みを書き忘れたクエリへの歯止めです。設定できるのはクエリごとの上限が1つ、ワークグループ全体の上限は複数。ワークグループはIAMリソースなので「このロールはapp-batchでしかクエリを実行できない」という制御も書けます。作成上限は1リージョン1アカウントあたり1000個です。

engine version 3 で壊れるクエリの実例と移行前テストの観点

自動アップグレードに任せる前に、どのクエリが壊れるかを知っておくべきでしょう。公式が破壊的変更として挙げるもののうち、既存バッチで踏みやすいのは次の5つです。

  • CONCAT が2引数以上を要求し、CONCAT(str)INVALID_FUNCTION_ARGUMENT で失敗する
  • GROUP BY のネスト列にダブルクォートが必須になり GROUP BY user.name が通らない
  • uuid() の戻り型が変わり、CTASやビューで CAST(uuid() AS VARCHAR) にしないと失敗する
  • 日時へのキャストで日付と時刻の間のハイフンが許されず、区切りは半角スペースが要る
  • charvarchar を混ぜた CONCAT|| 連結が型不一致で失敗する

いずれもエラーメッセージが明確なので、検証用ワークグループで既存クエリを一巡させれば洗い出せます。見落としやすいのは approx_percentile でしょう。エラーにならず、結果の数値だけが変わります。内部実装が qdigest から tdigest へ変わったためで、公式は近似値であることを理由にこの関数へ依存しないよう促しています。監視のしきい値をこれで出しているなら、移行を機に見直す対象です。

クエリ結果再利用の期限設定の既定60分と最大7日・効かなくなる条件

同じクエリを繰り返すダッシュボードやAPIでは、結果の再利用がそのままスキャン量の削減になります。公式の結果再利用のページによれば、期限を指定しない場合の既定は60分、指定できる最大は単位を問わず7日相当。再利用は同一ワークグループ内でのみ働き、対象は SELECTEXECUTE に限られます。

効かなくなる条件のほうが実務では効いてきます。フェデレーテッドカタログを参照するクエリ、複数のデータカタログにまたがるクエリ、参照テーブルが21個以上のクエリは対象外。Lake Formationで行・列のフィルタを掛けたテーブルを含む場合も使えません。クエリ文字列が100KB未満なら空白とコメントの差は無視され INNER JOINJOIN も同一視されますが、100KB超は完全一致が要ります。rand() のような非決定的な関数もキャッシュされません。

注意すべき性質が一つ。Athenaは指定期間が切れるまで元データの変更を見に行かないため、古い結果が返りえます。7日を指定した集計が、3日前に投入し直したデータを反映せず返る事態は普通に起こります。更新頻度の高いテーブルでは短く、日次で確定するテーブルでは長く、性質に合わせて期間を決めてください。

Athenaへの組み込みを進めてよい条件と別構成へ移すべき場面の線引き

ここは判断の章です。条件を示したうえで言い切ります。

同時実行クォータとキュー待ちから見て組み込みを避けるべき用途

画面表示のたびにAthenaを叩くアプリケーションは作りません。理由は3つあり、いずれもクォータに由来します。Active DMLのクォータに実行中とキュー待ちの両方が数えられるため、同時アクセスが増えた瞬間に TooManyRequestsException が返ること。開始から結果取得まで最低でも数百ミリ秒かかり、調整では縮まないこと。クエリ文字列の上限262144バイトが調整不可であることです。

秒単位の応答が要る参照系は、Athenaで集計した結果をRDSやDynamoDBへ書き出し、アプリはそちらを読む構成にします。Athenaは集計を作る側に置き、読ませる側には置かない。同時実行の心配が集計ジョブの本数だけに閉じます。パーティション数にも上限があり、Glueのテーブルが1000万パーティションを持てても、Athenaが1回のスキャンで読めるのは100万までです。時刻を分単位で切る設計は、数年分を溜めた時点でこの壁に当たります。

Athenaのまま進めてよい条件とRedshift Serverlessへ移すべき場面

Athenaのまま進めてよいのは、次の条件が揃うときです。クエリの発生源が人またはバッチで、同時実行が数十本に収まること。1件あたり数秒から数十秒の応答を許容できること。データがS3にあり、更新が要るならIcebergのv2テーブルで表現できること。この3つを満たすなら、クラスターを持たない分だけAthenaが有利になります。

移すべきなのは、同じ集計を高頻度で繰り返し、応答時間に責任を持つ場面です。結果再利用でしのげる範囲を超え、同時実行クォータの引き上げ申請を繰り返すようになったら構成の限界を示すサインでしょう。Redshift Serverlessの仕組みと移行判断で扱うような、ワークロードに応じて容量が伸びる仕組みのほうが素直に収まります。月に数回の調査のためにDWHを常設するのは、逆に過剰です。

迷いやすいのは中間。日次のバッチ集計とアドホック調査が同居する段階です。ここはAthenaのまま、ワークグループを2つに割るのが答えになります。設計と実装をまとめて外に出すなら、データ分析基盤の構築とMLOpsの支援でS3からAthena、BIまでの構成を引き受けています。

よくある質問

Athenaの組み込みで検索される質問のうち、公式ドキュメントに答えがあるものを5つ挙げます。

AthenaのクエリはSDKから同期的に実行できますか?

AWS SDKのAthenaクライアントに、投げて結果まで返す同期APIはありません。start_query_execution が返すのはクエリの実行IDだけで、get_query_execution で状態を確認し、SUCCEEDED を待ってから get_query_results で取得します。同期的な書き味が要るならJDBC/ODBCドライバがポーリングを隠してくれますが、サーバー側の待ち時間が消えるわけではありません。

フェデレーテッドクエリは追加費用がかかりますか?

追加費用が発生する仕組みです。AWS公式の料金ページは、S3以外のデータソースへのクエリについて、横断で集計したスキャン量をメガバイト単位に切り上げ、クエリごとに最小10MBを適用すると記載しています。さらにコネクタとして動くLambdaの実行料金が標準レートで別に発生します。数KBのマスタを引くだけでも10MB分が課金されるため、小さく高頻度な参照をこれで回す構成は費用が合いません。

Icebergテーブルのタイムトラベルはどこまで遡れますか?

遡れる範囲はスナップショットが残っている期間に等しく、固定の日数制限ではありません。VACUUM を実行すると古いスナップショットと不要ファイルが消え、その時点より前へは戻れなくなります。構文はengine version 3で FOR TIMESTAMP AS OFFOR VERSION AS OF に変わり、旧来の FOR SYSTEM_TIME AS OF はエラーです。どこまで戻す必要があるかを決めてから VACUUM の間隔を設計してください。

engine version 3 へのアップグレードで何が壊れますか?

既存クエリで踏みやすいのは、CONCAT の引数が2個以上必須になった点、GROUP BY のネスト列にダブルクォートが要る点、uuid() をCTASで直接使えなくなった点、日時キャストでハイフン区切りが許されなくなった点です。エラーにならず結果だけ変わる approx_percentile にも注意が要ります。検証用ワークグループを新バージョンで作り、既存クエリを一巡させてから本番を上げてください。

クエリ結果の再利用が効かないのはどんなときですか?

フェデレーテッドカタログを参照するクエリ、複数のデータカタログにまたがるクエリ、参照テーブルが21個以上のクエリ、Lake Formationで行や列のフィルタが掛かったテーブルを含むクエリでは働きません。クエリ側で出力先設定を上書きした場合や、rand() のような非決定的な関数を含む場合も対象外です。期限を指定しなければ既定60分、上限は7日相当で、期限内は古い結果が返ります。

関連記事

資料請求

RELATED POSTS 関連記事