リアルタイムでのデータ処理や分析が不可欠な現代のアプリケーション開発において、大量のデータを低遅延かつ高耐久に集約するシステムの構築は非常に重要です。本記事では、AWSが提供するフルマネージドなリアルタイムストリーミングデータサービスである Amazon Kinesis Data Streams (KDS) について、開発者が知っておくべきコアコンセプト、プロデューサーとコンシューマーの開発手法、そしてセキュリティやスケーリングのベストプラクティスを解説します。
1. Amazon Kinesis Data Streams(KDS)の基本アーキテクチャ
Amazon Kinesis Data Streamsは、数十万のソースからギガバイト単位のデータをリアルタイムで継続的に収集・保存できる、可用性と耐久性に優れたサーバーレスデータストリーミングサービスです。
KDSを構築する上で、以下の基本用語を理解することが開発の第一歩となります。
- シャード (Shard): ストリームのベースとなる容量単位です。1つのシャードは、書き込み時に最大 1 MB/秒 または 1,000レコード/秒、読み込み時に最大 2 MB/秒 の処理能力を提供します。
- データレコード (Data Record): ストリーム内に保存される最小のデータ単位です。レコードは、KDSから自動的に割り当てられるシーケンス番号 (Sequence Number)、データをシャードに分散させるためのパーティションキー (Partition Key)、および暗号化されていないバイト配列である最大 1 MB のデータBlob で構成されます。
- プロデューサー (Producer): KDSにレコードを送信する(インジェストする)アプリケーションやサービスです。
- コンシューマー (Consumer): KDSからレコードを取得してリアルタイムに処理を行うアプリケーション(Amazon Kinesis Data Streams アプリケーション)やAWSサービスです。
2. プロデューサーの開発:データ書き込みの最適化
KDSにデータを送信する方法として、開発者には主に以下の3つのアプローチが提供されています。
① Amazon Kinesis Data Streams API(AWS SDK)
最も直接的な方法で、PutRecord(単一のレコード送信) または PutRecords(複数のレコードをまとめて送信) のAPI操作を利用します。PutRecords はバッチ処理に対応しているため、HTTPコネクションのオーバーヘッドを大幅に削減し、送信パフォーマンスを向上させることができます。
② Amazon Kinesis Producer Library (KPL)
高いスループットを最小限のクライアントリソースで実現するための高度に最適化されたライブラリです。KPLはC++のバックエンドデーモンを利用した非同期アーキテクチャを採用しており、以下のバッチ処理を自動的に実行します。
- Aggregation(集約): 複数の小さなユーザーレコードをシリアライズし、単一のKDSレコードに詰め込むことで、PUTペイロード料金を削減し、1 MBのシャード制限を最大限に活用します。
- Collection(収集): 異なるパーティションキーや宛先シャードを持つ複数のKDSレコードをパッケージ化し、一度の
PutRecordsAPIコールで送信します。
③ Amazon Kinesis Agent
Linuxなどのサーバー環境にインストールして使用する事前構築済みのJavaアプリケーションです。ローカルのログファイルを監視し、ファイルのローテーションや再試行を処理しながらデータを継続的にストリームに送信します。また、送信前にマルチラインログを一行に変換するなどのデータ事前処理機能も備えています。
注意(サポート終了に関する重要なお知らせ) AWSは 2026年1月30日をもって、従来のKinesis Client Library (KCL) 1.x および Kinesis Producer Library (KPL) 0.x のサポート終了を発表しています。既存のアプリケーションを運用している場合は、最新の KPL 1.x および KCL 3.x への早急な移行を計画してください。
3. コンシューマーの開発:低遅延処理とスケーリング
インジェストされたデータを処理するコンシューマーの開発には、サーバーレスなアプローチとカスタム実装のアプローチがあります。
① AWS Lambda によるサーバーレス処理
AWS Lambdaは、KDSのイベントソースマッピング(ESM)を介してネイティブに統合されています。Lambdaサービスはバックグラウンドで各シャードを定期的にポーリングし、新しいレコードを検知するとLambda関数を同期的に呼び出します。
- Parallelization Factor(並列化要素): デフォルトでは1シャードにつき1つのLambda関数インスタンスしか実行されませんが、このパラメータを 1〜10 に設定することで、レコードの処理順序を保証したまま、シャードあたり最大10個の並列Lambda呼び出しを可能にします。
- 部分バッチ失敗報告 (batchItemFailures): バッチの一部に処理失敗のレコード(いわゆる毒薬レコード)が含まれている場合、失敗したレコードのみを再試行の対象として報告し、成功したものは正しくコミット(チェックポイント)することで再試行コストと遅延を最小化できます。
② Kinesis Client Library (KCL) によるカスタム処理
サーバー管理が伴う大規模かつ複雑な分散ストリーム処理アプリケーション(Amazon EC2、Amazon ECS、EKS、Fargateなどで稼働)を開発する場合、KCLが最適な選択肢となります。
KCLは、DynamoDBの「リース表(Lease Table)」を使用してシャードの割り当てや進捗状況(チェックポイント)をメタデータとして協調管理します。これにより、複数コンシューマーワーカー間での自動負荷分散や障害発生時のフェイルオーバーが実現します。
- シングルテーブルフォーマット (KCL 3.5以降): KCL 3.xでは、従来作成されていた「リース表」「ワーカーメトリクス表」「コーディネーターステート表」の3つのDynamoDBテーブルを単一のリース表に統合することができます。これによりDynamoDBのアカウント内テーブル制限を回避し、効率的なメタデータ管理が可能となります。
4. キャパシティ管理:ワークロードに合わせたモードの使い分け
KDSには、ユースケースやコスト構造に応じて選択可能な 3つの容量モード があります。
| 容量モード | スケーリングの管理 | 最適なユースケース |
|---|---|---|
| Provisioned(プロビジョンド) | 手動でシャード数を指定。UpdateShardCount APIなどでスケーリング。 | トラフィックが予測可能で安定したワークロード。 |
| On-Demand Standard(オンデマンド・スタンダード) | 自動スケール。データ書き込み、読み込み、ストレージの量に応じてGB単位で課金。 | トラフィックが予測不可能でスパイクの多いワークロード。 |
| On-Demand Advantage(オンデマンド・アドバンテージ) | アカウントレベルの設定で自動有効化。 aggregate 10 MiB/s以上の書き込みを伴う場合にスループット課金が大幅値引きされる。 | 10 MiB/s以上の高トラフィックパイプライン、または1リージョン内に多数のオンデマンドストリームを稼働させる場合。 |
- ホットシャードへの対応(プロビジョンドモードの場合): 特定のパーティションキーへの偏りによって特定のシャードが過負荷(ホットシャード)になった場合、リシャーディング操作(シャードの分割: Split または隣接するシャードの結合: Merge)をAPI経由で手動で実行してストリームを最適化します。
5. セキュリティとコンプライアンスの担保
KDSは安全に機密データを送信・保管するために必要な、エンタープライズレベルのセキュリティ機能を標準装備しています。
- データ転送中の保護 (Data in Transit): KDS APIへの接続はすべて TLS 1.2 以上 のプロトコルによるHTTPS通信が強制されます。
- データ保管時の保護 (Data at Rest): AWS KMSを利用した サーバーサイド暗号化 (SSE-KMS) がサポートされています。AWS管理のキー(
aws/kinesis)も選択可能ですが、監査、ローテーション、およびきめ細かなキーポリシーを制御するために、カスタマー管理のKMSキー (CMK) を作成して適用することがベストプラクティスです。 - プライベート接続 (AWS PrivateLink): インターネットを経由せず、自身のVPC(Virtual Private Cloud)内からプライベートIPアドレスを使用してKDS APIを安全に呼び出すための インターフェイスVPCエンドポイント を作成可能です。
- リソースベースポリシーとクロスアカウントアクセス: KDSはリソースベースポリシーをサポートしており、異なるAWSアカウント間でスループット共有を設定したり、他アカウントのLambda関数を直接トリガーするセキュアなパイプラインを簡単に構成できます。
まとめ
Amazon Kinesis Data Streamsは、インフラの管理負荷を排除し、マイクロサービスアーキテクチャやイベント駆動システムにおける強力なリアルタイム基盤として機能します。KPL/KCLといった堅牢なライブラリやLambdaのParallelization Factorなどの制御パラメータを活用し、要件に最適化された高性能なパイプラインを設計しましょう。
引用先(AWS公式ドキュメントおよび仕様)
- Amazon Kinesis Data Streams Terminology and concepts https://docs.aws.amazon.com/streams/latest/dev/key-concepts.html
- Develop producers using the Amazon Kinesis Producer Library (KPL) https://docs.aws.amazon.com/streams/latest/dev/developing-producers-with-kpl.html
- Choose the right mode to stream in – Amazon Kinesis Data Streams https://docs.aws.amazon.com/streams/latest/dev/how-do-i-size-a-stream.html
- DynamoDB metadata tables and load balancing in KCL https://docs.aws.amazon.com/streams/latest/dev/kcl-dynamoDB.html
- Single table format for KCL – Amazon Kinesis Data Streams https://docs.aws.amazon.com/streams/latest/dev/kcl-single-table-format.html
- What is server-side encryption for Kinesis Data Streams? https://docs.aws.amazon.com/streams/latest/dev/what-is-sse.html
- Use Amazon Kinesis Data Streams with interface VPC endpoints https://docs.aws.amazon.com/streams/latest/dev/vpc.html
- Controlling access to Amazon Kinesis Data Streams resources using IAM https://docs.aws.amazon.com/streams/latest/dev/controlling-access.html
- Quotas and limits – Amazon Kinesis Data Streams https://docs.aws.amazon.com/streams/latest/dev/service-sizes-and-limits.html
