Kinesis Data Streams の新機能 「Streaming tables」 で S3 Tables への直接配信を試してみた
はじめに
2026年8月28日、Amazon Kinesis Data Streams の Streaming tables が発表されました。コンシューマを書かずに、ストリームに流れるデータを S3 Tables 上の Apache Iceberg テーブルへ配信できる機能です。翌29日には汎用 S3 バケットへの配信も発表され、KDS 単体で選べる配信先は、S3 Tables と汎用 S3 バケットの2種類になりました。
本記事では、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.json と bad-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 で分析するなら、試してみてください。







