Kinesis Data Streams の新機能 「Streaming tables」 で S3 Tables への直接配信を試してみた

Kinesis Data Streams の新機能 「Streaming tables」 で S3 Tables への直接配信を試してみた

Amazon Kinesis Data Streams が Streaming tables を発表し、ストリームのデータを S3 Tables 上の Apache Iceberg テーブルへ直接配信できるようになりました。同じストリームに S3 Tables 宛先と汎用 S3 バケット宛先の配信を1本ずつ作成し、Athena でクエリした際の違いを確かめました。
2026.09.03

はじめに

2026年8月28日、Amazon Kinesis Data Streams の Streaming tables が発表されました。コンシューマを書かずに、ストリームに流れるデータを S3 Tables 上の Apache Iceberg テーブルへ配信できる機能です。翌29日には汎用 S3 バケットへの配信も発表され、KDS 単体で選べる配信先は、S3 Tables と汎用 S3 バケットの2種類になりました。

https://aws.amazon.com/about-aws/whats-new/2026/08/kinesis/data-delivery-s3-tables

https://dev.classmethod.jp/articles/kinesis-data-streams-s3-general-purpose-delivery/

本記事では、1つのストリームに S3 Tables 宛先と汎用 S3 バケット宛先の配信を1本ずつ作成しました。同じレコードが両方の宛先へどう届いたかを、参照先だけを差し替えた同じ SQL を Athena で実行して比べました。

検証内容

Streaming tables の前提

Streaming tables と汎用 S3 配信は、どちらも On-Demand Standard または On-Demand Advantage のストリームが対象です。Provisioned モードのストリームでは使えません(開発者ガイド)。1ストリームに立てられる配信は、S3 Tables 向け1本と汎用 S3 向け1本の計2本までで、緩和できません(同じクォータ表)。

S3 Tables への配信には、汎用 S3 バケットへの配信にはない追加要件があります。

S3 Tables への配信では、AWS Glue Schema Registry に登録したスキーマの指定が必須です。配信できなかったレコードのエラー情報の出力先となる、デッドレターキュー用の S3 バケットも必須です(開発者ガイド)。パーティションに使う列は、Iceberg テーブルで timestamptz 型である必要があります。JSON Schema では string 型に format: date-time を付けて定義します(開発者ガイド)。クロスアカウント配信とクロスリージョン配信はできないため、ストリーム、テーブルバケット、Glue Schema Registry を同一アカウント・同一リージョンに置きます(開発者ガイド)。

検証環境

項目
検証日 2026-09-02
AWS CLI aws-cli/2.36.36
リージョン ap-northeast-1
ストリームのキャパシティモード On-Demand Standard
宛先の暗号化 SSE-S3(既定)

関連リソースの準備

配信の宛先と入力元になるリソースを作成しました。

# 入力元の ODS ストリーム
aws kinesis create-stream --stream-name kdsst-0902-stream \
  --stream-mode-details StreamMode=ON_DEMAND

# S3 Tables 配信の宛先になるテーブルバケット
aws s3tables create-table-bucket --name kdsst-0902-tb

# 汎用 S3 配信の宛先とデッドレターキューを兼ねる汎用 S3 バケット
aws s3api create-bucket --bucket kdsst-0902-123456789012 \
  --create-bucket-configuration LocationConstraint=ap-northeast-1

次に、S3 Tables 配信で必須になる JSON Schema を event-schema.json として用意しました。

{
  "$schema": "http://json-schema.org/draft-07/schema#",
  "title": "event",
  "type": "object",
  "properties": {
    "event_time": { "type": "string", "format": "date-time" },
    "user_id": { "type": "string" },
    "action": { "type": "string" },
    "region": { "type": "string" },
    "latency_ms": { "type": "integer" },
    "payload": { "type": "string" }
  },
  "required": ["event_time", "user_id", "action"]
}

required に挙げた項目は、Iceberg テーブルで NOT NULL の列になります。パーティションに使う event_time が、上に書いた timestamptz の列に当たります。

このファイルを Glue Schema Registry へ登録しました。レジストリを作成してから、スキーマを登録しました。

aws glue create-registry --registry-name kdsst-0902-registry

aws glue create-schema \
  --registry-id RegistryName=kdsst-0902-registry \
  --schema-name events \
  --data-format JSON \
  --compatibility NONE \
  --schema-definition file://event-schema.json

データ形式は JSON、互換性は NONE を指定しました。作成結果に含まれるスキーマの ARN は、続く配信の作成で指定しました。

Kinesis Data Streams が配信処理で引き受けるサービス実行ロールを作成しました。信頼ポリシーは trust-policy.json として用意しました。

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Principal": { "Service": "kinesis.amazonaws.com" },
      "Action": "sts:AssumeRole",
      "Condition": {
        "StringEquals": { "aws:SourceAccount": "123456789012" },
        "ArnLike": { "aws:SourceArn": "arn:aws:kinesis:ap-northeast-1:123456789012:channel/*" }
      }
    }
  ]
}

権限ポリシーは permission-policy.json にまとめました。本検証で設定した権限は、宛先のテーブルバケットに対する S3 Tables のアクション、スキーマを読むための Glue Schema Registry のアクション、汎用 S3 バケットへの書き込みの3種類です。後者のバケットは、汎用 S3 配信の宛先とデッドレターキューを兼ねています。S3 Tables のアクションには、開発者ガイドのポリシー例に挙がっているものを、CMK を使わない構成でもすべて付与しました。

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "S3TablesAccess",
      "Effect": "Allow",
      "Action": [
        "s3tables:GetTable",
        "s3tables:GetTableBucket",
        "s3tables:GetTableMetadataLocation",
        "s3tables:UpdateTableMetadataLocation",
        "s3tables:CreateTable",
        "s3tables:CreateNamespace",
        "s3tables:PutTableData",
        "s3tables:GetTableData",
        "s3tables:TagResource",
        "s3tables:PutTableRecordExpirationConfiguration",
        "s3tables:PutTableEncryption"
      ],
      "Resource": [
        "arn:aws:s3tables:ap-northeast-1:123456789012:bucket/kdsst-0902-tb",
        "arn:aws:s3tables:ap-northeast-1:123456789012:bucket/kdsst-0902-tb/table/*"
      ]
    },
    {
      "Sid": "GlueSchemaRegistryAccess",
      "Effect": "Allow",
      "Action": [
        "glue:GetSchemaVersion"
      ],
      "Resource": [
        "arn:aws:glue:ap-northeast-1:123456789012:registry/kdsst-0902-registry",
        "arn:aws:glue:ap-northeast-1:123456789012:schema/kdsst-0902-registry/events"
      ]
    },
    {
      "Sid": "DeliveryBucketList",
      "Effect": "Allow",
      "Action": [
        "s3:ListBucket",
        "s3:ListBucketMultipartUploads"
      ],
      "Resource": [
        "arn:aws:s3:::kdsst-0902-123456789012",
        "arn:aws:s3:::kdsst-0902-123456789012/*"
      ]
    },
    {
      "Sid": "DeliveryBucketWrite",
      "Effect": "Allow",
      "Action": [
        "s3:PutObject",
        "s3:CreateMultipartUpload",
        "s3:UploadPart",
        "s3:CompleteMultipartUpload",
        "s3:ListMultipartUploads",
        "s3:ListMultipartUploadParts"
      ],
      "Resource": [
        "arn:aws:s3:::kdsst-0902-123456789012/*"
      ]
    }
  ]
}

ロールを作成し、このポリシーを付与しました。

aws iam create-role --role-name kdsst-0902-delivery-role \
  --assume-role-policy-document file://trust-policy.json

aws iam put-role-policy --role-name kdsst-0902-delivery-role \
  --policy-name kdsst-0902-delivery-policy \
  --policy-document file://permission-policy.json

配信の作成

S3 Tables 宛先の create-channel に指定した JSON です。create-channel-tables.json として保存しました。

{
  "ChannelName": "kdsst-0902-tables",
  "ServiceExecutionRoleARN": "arn:aws:iam::123456789012:role/kdsst-0902-delivery-role",
  "StreamConfigurationList": [
    {
      "StreamARN": "arn:aws:kinesis:ap-northeast-1:123456789012:stream/kdsst-0902-stream",
      "RecordConfiguration": {
        "RecordFormatType": "JSON",
        "GSRSchemaARN": "arn:aws:glue:ap-northeast-1:123456789012:schema/kdsst-0902-registry/events"
      }
    }
  ],
  "S3TablesDestinationConfiguration": {
    "DataFreshnessInSeconds": 300,
    "DeadLetterQueueS3Configuration": {
      "BucketARN": "arn:aws:s3:::kdsst-0902-123456789012",
      "ExpectedBucketOwner": "123456789012",
      "ErrorOutputPrefix": "dlq-tables/"
    },
    "S3TablesConfigurationList": [
      {
        "TableBucketARN": "arn:aws:s3tables:ap-northeast-1:123456789012:bucket/kdsst-0902-tb",
        "Namespace": "kdsdemo",
        "TableName": "events",
        "CompressionType": "ZSTD",
        "PartitionSpec": {
          "PartitionFields": [
            { "Transform": "TIME_HOUR", "SourceName": "event_time" }
          ]
        }
      }
    ]
  }
}

RecordFormatType に JSON を指定して GSRSchemaARN を渡すと、Glue Schema Registry のシリアライザでエンコードしていない JSON レコードを投入できます。Glue Schema Registry のシリアライザでレコードにスキーマIDを埋め込む GSR_JSON も選べます(開発者ガイド)。今回は JSON を使いました。

続いて、汎用 S3 宛先の create-channel に指定した JSON です。create-channel-s3.json として保存しました。

{
  "ChannelName": "kdsst-0902-s3",
  "ServiceExecutionRoleARN": "arn:aws:iam::123456789012:role/kdsst-0902-delivery-role",
  "StreamConfigurationList": [
    {
      "StreamARN": "arn:aws:kinesis:ap-northeast-1:123456789012:stream/kdsst-0902-stream",
      "RecordConfiguration": { "RecordFormatType": "JSON" }
    }
  ],
  "S3DestinationConfiguration": {
    "DataFreshnessInSeconds": 300,
    "DeadLetterQueueS3Configuration": {
      "BucketARN": "arn:aws:s3:::kdsst-0902-123456789012",
      "ExpectedBucketOwner": "123456789012",
      "ErrorOutputPrefix": "dlq-s3/"
    },
    "StorageConfiguration": {
      "BucketARN": "arn:aws:s3:::kdsst-0902-123456789012",
      "ExpectedBucketOwner": "123456789012",
      "OutputKeyTemplate": "data/!{yyyy}/!{MM}/!{dd}/!{HH}/!{channel-name}!{extension}",
      "StorageClass": "STANDARD",
      "CompressionType": "GZIP"
    }
  }
}

それぞれの入力ファイルを create-channel に指定して、配信を作成しました。

aws kinesis create-channel --cli-input-json file://create-channel-tables.json

aws kinesis create-channel --cli-input-json file://create-channel-s3.json

宛先の指定は S3TablesDestinationConfiguration と S3DestinationConfiguration に分かれています。2本の create-channel を 04:12:25Z と 04:12:46Z に実行し、04:12:57Z 時点でどちらも ACTIVE になっていました。

宛先のテーブルは、今回の配信作成時に Streaming tables によって自動作成されました。既存のテーブルへは配信できません(開発者ガイド)。

自動作成されたテーブルの情報
{
    "name": "events",
    "type": "customer",
    "tableARN": "arn:aws:s3tables:ap-northeast-1:123456789012:bucket/kdsst-0902-tb/table/<table-id>",
    "namespace": [
        "kdsdemo"
    ],
    "metadataLocation": "s3://<table-id-prefix>-<suffix>--table-s3/metadata/00001-<metadata-id>.metadata.json",
    "warehouseLocation": "s3://<table-id-prefix>-<suffix>--table-s3",
    "createdBy": "123456789012",
    "managedByService": "kinesis.amazonaws.com",
    "ownerAccountId": "123456789012",
    "format": "ICEBERG"
}

自動作成されたテーブルは、managedByService が kinesis.amazonaws.com のサービス管理テーブルです。スキーマやテーブルプロパティを外から書き換える対象ではありません(開発者ガイド)。

Athena では、JSON Schema の列が次の型で参照できました。

JSON Schema での定義 Athena で見える型
event_time string + format: date-time timestamp with time zone
latency_ms integer integer
user_id string varchar
action string varchar
region string varchar
payload string varchar

event_time だけが timestamp with time zone 型になり、ほかの列は JSON Schema の型に対応する Athena の型として参照できました。

データ投入

2本の配信が動いている状態で、ストリームにレコードを投入しました。

項目
スキーマに沿ったレコード 53,000 件
スキーマに沿わないレコード 2 件
1件あたりの大きさ 約 1,033 バイト
PutRecords の失敗 0 件

payload 列には、同じ文字を 900 文字並べた値を入れました。gzip が効きやすい内容のため、以降のサイズとスキャン量の比較結果は、汎用 S3 配信に有利な条件での測定値です。

スキーマに沿わないレコードとして投入したのは、action を欠いた1件と、latency_ms に文字列を入れた1件です。bad-record-1.jsonbad-record-2.json に分けて保存しました。

{"event_time":"2026-09-02T04:15:00.000Z","user_id":"user-bad1","region":"ap-northeast-1","latency_ms":10,"payload":"missing-action"}
{"event_time":"2026-09-02T04:15:01.000Z","user_id":"user-bad2","action":"view","region":"ap-northeast-1","latency_ms":"not-a-number","payload":"bad-type"}

この2件は PutRecord で個別に投入しました。生の JSON を渡すため、--cli-binary-format に raw-in-base64-out を指定しました。

aws kinesis put-record --stream-name kdsst-0902-stream \
  --partition-key bad1 \
  --cli-binary-format raw-in-base64-out \
  --data file://bad-record-1.json

DataFreshnessInSeconds は両方 300 に設定しましたが、最初のオブジェクトと Iceberg テーブルの作成を確認できたのは投入から約13分後でした。初回は宛先の準備を含むため、この値どおりには届きません。

投入したデータの確認

配信の処理量は CloudWatch メトリクスで確認しました。AWS/Kinesis 名前空間で、2026-09-02 04:10Z〜05:15Z を対象に統計値 Sum で集計した、DeliveryToIceberg と DeliveryToS3 の接頭辞が付いたメトリクスです。

メトリクス Streaming tables 汎用 S3 配信
BytesIn 55,239,670 55,239,670
BytesProcessed 54,762,662 54,762,662
BytesOut 264,801 799,726
RecordCount 値なし(今回の取得) 53,002
SuccessfulRecordCount 値なし(今回の取得) 53,002
FailedRecordCount 値なし(今回の取得) 0
DLQDeliverySuccess 2 値なし(今回の取得)

入力側のバイト数は、2本の配信で完全に一致しました。この結果から、2本の配信が同じレコードをそれぞれ独立に読み取ったことを確認できました。出力量は Streaming tables 側が 264,801 バイトで、汎用 S3 配信の 799,726 バイトに対して約3分の1でした。前者は Parquet と ZSTD、後者は GZIP の JSON です。汎用 S3 配信の実オブジェクトは3個で、合計サイズ 799,726 バイトが BytesOut と一致しました。

この取得期間では、Streaming tables 側のレコード件数系メトリクスには値が返らず、Bytes 系と DLQDeliverySuccess の値だけを確認できました。この1回の観測範囲から、メトリクスの提供可否を一般化するものではありません。

カタログは s3tablescatalog/<テーブルバケット名>、データベースは namespace、テーブルは配信で指定したテーブル名として参照します。

SELECT action, count(*) AS n
FROM "s3tablescatalog/kdsst-0902-tb"."kdsdemo"."events"
GROUP BY action
ORDER BY action

今回のアカウントでは、分析サービスとの統合があらかじめ有効化されていたため、Lake Formation で追加の権限を付けずに参照できました。s3tablescatalog という Glue の連携カタログが、アカウント内のテーブルバケット全体を対象として既に存在していました。既定権限も IAM_ALLOWED_PRINCIPALS に ALL でした。この連携が入っていないアカウントでは、先に S3 Tables 側で分析サービスとの統合を有効化する必要があります(ユーザーガイド)。

続いて、Athena のスキャン量を比べました。汎用 S3 配信のほうは、配信された GZIP の JSON を読む外部テーブルを作成しました。

CREATE EXTERNAL TABLE IF NOT EXISTS kdsst_0902.events_json (
  event_time string,
  user_id string,
  action string,
  region string,
  latency_ms bigint,
  payload string
)
ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe'
LOCATION 's3://kdsst-0902-123456789012/data/'

実行した SQL は FROM の参照先だけが異なり、選択する列と集計は両方で同じでした。

問い合わせ Streaming tables 汎用 S3 配信
action ごとの件数集計 8,367 バイト 799,726 バイト
全件数の取得 0 バイト 799,726 バイト

集計では、約96分の1のスキャン量で済みました。Parquet は列単位で保存されているので action 列だけを読めばよく、JSON 側はオブジェクト全体を読む必要があります。全件数の取得は Streaming tables 側がスキャン 0 バイトでした。Iceberg のメタデータだけで応答できたためです。

action ごとの件数集計の結果を並べます。

action Streaming tables 汎用 S3 配信
click 10,600 10,600
logout 10,600 10,600
purchase 10,600 10,600
signup 10,600 10,600
view 10,600 10,601
(NULL) 1
合計 53,000 53,002

スキーマに沿わない2件の扱いが、宛先で分かれました。Streaming tables 側では、スキーマ検証に失敗した2件がデッドレターキューへ送られ、テーブルには入りませんでした。汎用 S3 配信ではレコードが検証されずに格納され、Athena の結果では、action を欠いた1件が action が NULL の行として現れました。latency_ms に文字列を入れた1件は、action が view の行として集計されました。

デッドレターキューのオブジェクトは、dlq-tables/DESERIALIZATION_ERROR/2026/09/02/04/ の下に置かれました。指定した ErrorOutputPrefix の後にエラー種別のディレクトリが挟まり、その下に時刻のディレクトリが続きます。

デッドレターキューのオブジェクトの中身
{"approximateArrivalTimestamp":1788322490093,"streamArn":"arn:aws:kinesis:ap-northeast-1:123456789012:stream/kdsst-0902-stream","shardId":"shardId-000000000001","sequenceNumber":"49677952019003285752333393925921343430349053876775157778","errorCode":"Iceberg.MissingColumnWithinRecord","errorMessage":"Data Channel is unable to deliver to the table since a required column within the schema is missing within the record."}
{"approximateArrivalTimestamp":1788322490551,"streamArn":"arn:aws:kinesis:ap-northeast-1:123456789012:stream/kdsst-0902-stream","shardId":"shardId-000000000003","sequenceNumber":"49677952019047887242730455172492139211962632343367188530","errorCode":"Iceberg.MalformedFieldWithinRecord","errorMessage":"Data Channel is unable to convert column data in your record to the column type specified within the schema."}

レコード本体は含まれていませんでした。今回のオブジェクトには、ストリームの ARN、シャードID、シーケンス番号、エラーの種別だけが、1行1レコードの JSON 形式で書き出されていました。再投入するには、シーケンス番号を頼りにストリームから読み直す必要があります。

撤去

配信は削除するまで課金が続くため、検証を終えた後に配信2本を先に削除しました。テーブルとネームスペースは配信が自動で作成するため、テーブルバケットを削除する前に個別に削除しました。今回の撤去はこの順で実行しました。

aws kinesis delete-channel --channel-arn <channel-arn>
aws kinesis delete-stream --stream-name kdsst-0902-stream --enforce-consumer-deletion

aws s3tables delete-table --table-bucket-arn <table-bucket-arn> --namespace kdsdemo --name events
aws s3tables delete-namespace --table-bucket-arn <table-bucket-arn> --namespace kdsdemo
aws s3tables delete-table-bucket --table-bucket-arn <table-bucket-arn>

aws s3 rm s3://kdsst-0902-<account-id> --recursive
aws s3api delete-bucket --bucket kdsst-0902-<account-id>

aws glue delete-schema --schema-id SchemaArn=<schema-arn>
aws glue delete-registry --registry-id RegistryName=kdsst-0902-registry
aws glue delete-table --database-name kdsst_0902 --name events_json

aws iam delete-role-policy --role-name kdsst-0902-delivery-role --policy-name kdsst-0902-delivery-policy
aws iam delete-role --role-name kdsst-0902-delivery-role

料金

配信の料金は宛先ごとに単価が違います。On-Demand Standard、US East (Ohio) の単価を、月 15,000 GB(500 GB/日 × 30 日)に揃えて計算しました。料金ページの計算例は宛先ごとに前提のデータ量が違うためです。

宛先 単価 15,000 GB/月 の料金
Streaming tables $0.035/GB $525.00
汎用 S3 配信 $0.0275/GB $412.50

Streaming tables の $0.035/GB は、Amazon Data Firehose 経由で Iceberg テーブルへ配信する場合より安く設定されています。Firehose の単価は、Kinesis Data Streams をソースにした Iceberg 宛先で $0.045/GB でした(Amazon Data Firehose の料金の計算例)。

課金の対象は、配信に成功したデータ量です。Streaming tables はその圧縮前のサイズで数えるので、Parquet に変換して小さくなっても配信の料金は下がりません。小さくなった分が効くのは Athena のスキャン課金($5/TB)で、こちらは読む量に比例します。

これとは別に、ストリーム本体の料金がかかります。On-Demand Standard の課金対象は、ストリーム時間、データ入力、データ取得、データ保持期間です。

まとめ

Streaming tables により、Kinesis Data Streams のデータを、自作のコンシューマを挟まずに S3 Tables の Iceberg テーブルへ配信できました。Parquet と Iceberg のメタデータとして届くため、Athena は必要な列だけを読み、全件数はメタデータだけで返します。

自作のコンシューマや別サービスを経由して Iceberg テーブルへ書いている構成は、配信先のテーブルを作り直せるなら Streaming tables に置き換えられます。Amazon Data Firehose を挟む構成と比べても単価は安く、必要なのはチャネル1本です。Parquet と Iceberg のメタデータを活かせる Athena で分析するなら、試してみてください。

この記事をシェアする

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

関連記事