0Pricing
Cloud & IT Cert Prep · レッスン

リアルタイムイベント処理向けKinesis Data Streams

Kinesis Data Streamsで高スループットのイベントストリームを生成・消費し、スループットに応じてシャードを管理して、Lambdaをコンシューマーとして使用します。

「リアルタイムイベント処理向けKinesis Data Streams」はCoddyKit上の無料Cloud & IT Cert Prepレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはCloud & IT Cert Prep学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Cloud & IT Cert Prepコースには全4レッスンが含まれています。

Kinesis Data Streams の基本概念

Kinesis Data Streams (KDS) は、耐久性があり、順序性を備えたリアルタイムデータストリーミングサービスです。データは 1 つ以上の シャードで構成される ストリームに整理されます。各シャードは、データレコードの順序付きシーケンスです。プロデューサーは、レコードを受け取るシャードを決定する パーティションキーを使用して、シャードにレコードを書き込みます。コンシューマーはシャードからレコードを読み取り、各シャード内で到着した順序どおりに処理します。

シャードの容量とスループット制限

各シャードは、書き込みスループット 1 MB/s または 1,000 レコード/sと、読み取りスループット 2 MB/s(そのシャード上のすべての標準コンシューマーで共有)に対応します。ストリーム全体の容量は、シャード数に比例して増加します。次の式を使用します。shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2) プロデューサーが書き込み制限に達すると、ProvisionedThroughputExceededException エラーが発生します。この場合は、シャードを分割するか、パーティションキーをより均等に分散させて解決します。

# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed

aws kinesis create-stream \
  --stream-name iot-telemetry \
  --shard-count 5

パーティションキーとデータ分散

パーティションキーは文字列であり、Kinesis はこれをハッシュ(MD5)して、レコードを受け取るシャードを決定します。適切に選択したパーティションキーを使用すると、レコードがシャード間で均等に分散され、ホットシャードを防止できます。IoT ではデバイス ID を使用します。クリックストリームではセッション ID またはユーザー ID を使用します。値の種類が少ないキー(たとえば、値が 5 つしかない国名)は避けてください。このようなキーを使用すると、1 つのシャードに書き込みトラフィックが偏り、他のシャードがアイドル状態になるホットシャードが発生します。

import boto3, json

client = boto3.client('kinesis', region_name='us-east-1')

# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
    StreamName='iot-telemetry',
    Data=json.dumps(event),
    PartitionKey='sensor-42'  # high-cardinality -> even distribution
)

標準コンシューマーと拡張ファンアウト

標準コンシューマーは、ポーリングで GetRecords を使用し、シャードあたり 2 MB/s の読み取りスループットを共有します。1 つのシャードに 3 つのコンシューマーがあり、それぞれが 2 MB/s を必要とする場合、合計 2 MB/s を共有するためスロットリングされます。Enhanced Fan-Out (EFO) を使用すると、登録済みの各コンシューマーが、永続的な HTTP/2 プッシュ接続(SubscribeToShard)を介して専用の 2 MB/s の読み取りパイプを利用できます。EFO ではコンシューマー・シャード時間ごとの料金が発生しますが、読み取り競合を完全に解消できます。

# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --consumer-name real-time-analytics

# List registered consumers
aws kinesis list-stream-consumers \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry

Kinesis コンシューマーとしての Lambda

Lambda は Event Source Mapping を介して Kinesis Data Streams とネイティブに統合できます。Lambda はストリームをポーリングしてレコードをバッチで読み取り、関数を呼び出します。BatchSize(1~10,000 レコード)、StartingPosition(最も古いレコードの場合は TRIM_HORIZON、最新の場合は LATEST)、および失敗したバッチを分割する BisectBatchOnFunctionError を設定します。Parallelisation Factor(1~10)を使用すると、Lambda はシャードごとに複数の呼び出しを同時に実行でき、高速に流入するストリームに対応できます。

# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
  --function-name ProcessIoTEvents \
  --event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --batch-size 100 \
  --starting-position LATEST \
  --parallelization-factor 5 \
  --bisect-batch-on-function-error true \
  --destination-config '{
    "OnFailure": {
      "Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
    }
  }'

Kinesis Client Library (KCL)

Kinesis Client Library (KCL) は、堅牢な Java(または MultiLangDaemon を介した複数言語対応)Kinesis コンシューマーを構築するためのアプリケーションフレームワークです。KCL は、シャードの列挙、リース管理(ワーカーインスタンス間でのシャード分配)、DynamoDB への進捗のチェックポイント保存、シャードの分割や結合への適切な対応を処理します。各 KCL ワーカーは 1 つ以上のシャードを処理し、ワーカーのスケールアウトや障害発生時には、KCL がシャードを自動的に再分散します。本番環境のコンシューマーアプリケーションでは、SDK を直接ポーリングするよりも KCL が推奨されます。

# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
#   'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();

データ保持と再生

Kinesis Data Streams は、デフォルトで24 時間レコードを保存します(追加料金で 7 日間、または Long-Term Retention を使用して最長 365 日間に延長できます)。SQS とは異なり、読み取ったレコードは消費後に削除されません。保持期間が समाप्तするまで利用できます。これにより、複数のコンシューマーが同じレコードをそれぞれ独立して読み取ることができ、コンシューマーのチェックポイントを以前のシーケンス番号に戻すことで再生も可能になります。これは、バグ修正や新しいサービスへのバックフィルに非常に役立ちます。

# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
  --stream-name iot-telemetry \
  --retention-period-hours 168

# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
  --stream-name iot-telemetry \
  --shard-id shardId-000000000000 \
  --shard-iterator-type TRIM_HORIZON \
  --query 'ShardIterator' --output text)

aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100

オンデマンドキャパシティモード

Kinesis Data Streams は 2 つのキャパシティモードをサポートしています。プロビジョンドモードでは、シャード数を手動で管理し、シャード時間単位で料金を支払います。オンデマンドモードでは、Kinesis が受信スループットに基づいてシャード容量を自動的にスケーリングします(デフォルトで書き込み最大 200 MB/s、読み取り最大 400 MB/s)。料金は書き込みおよび取得したデータの GB 単位で発生します。シャードのスケーリングを管理せずに、変動するトラフィックや予測困難なトラフィックに対応したい場合は、オンデマンドモードが適しています。

# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
  --stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
  --stream-mode-details StreamMode=ON_DEMAND

シャード内の順序保証

Kinesis はシャード内の順序を保証します。同じパーティションキーを持つレコードは常に同じシャードに送られ、書き込まれた順序で読み取られます。ただし、シャード間の順序は保証されません。アプリケーションで全レコードのグローバルな順序が必要な場合は、単一のシャードを使用して(スループットは 1 MB/s に制限されます)対応するか、パーティションキーグループ内(たとえばデバイス単位)でのみ順序が必要になるように設計し直します。これは、厳密な重複排除と順序付けを提供する SQS FIFO との重要な試験上の違いです。

セキュリティ:暗号化と VPC

Kinesis Data Streams は、サーバー側の暗号化が有効な場合、AWS KMS(CMK または AWS 管理キー)を使用してレコードをサーバー側で保管時に暗号化します。転送中のデータはすべて TLS で暗号化されます。VPC で実行され、ストリームデータをパブリックインターネット経由で送信したくないアプリケーションでは、Kinesis 用の VPC Interface Endpoint(PrivateLink)を使用してください。これにより、トラフィックを AWS ネットワークバックボーン内に完全に留めることができ、規制対象のワークロードで重要です。

# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
  --stream-name iot-telemetry \
  --encryption-type KMS \
  --key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123

Kinesis Data Streams、SQS、Kafka の比較

SAA-C03 試験では、Kinesis Data Streams と代替サービスを比較できるようにしておきましょう。Kinesis と SQSの違いは、Kinesis がシャード内の順序を保持し、複数のコンシューマーが同じデータを読み取れるのに対し、SQS はメッセージを消費後に削除する点です。Kinesis と MSK (Kafka)の違いは、MSK がマネージド Apache Kafka であり、Kafka プロトコルとの互換性、高度なトピック設定、オンプレミスの Kafka からの移行が必要な場合に使用する点です。AWS ネイティブのストリーミングと Lambda、Firehose、Flink との緊密な統合には Kinesis を使用します。問題文で Kafka または Kafka 互換性の要件が明示されていない限り、Kinesis を選択してください。

理解度チェック

このレッスンで学んだ AWS Solutions Architect(SAA-C03)の概念について、理解度を確認します。

レッスンのまとめ

このレッスンでは、Kinesis Data Streams が、シャードに分割され、順序性と耐久性を備えたストリームを提供し、デフォルトで 24 時間の保持と再生機能を備えていること、Enhanced Fan-Out により各コンシューマーがシャードごとに専用の 2 MB/s を利用でき、読み取り競合を解消できること、そして オンデマンドモードが予測困難なトラフィックに合わせてシャードを自動的にスケーリングすることを学びました。次は、イベント駆動アーキテクチャにおけるコレオグラフィーとオーケストレーションのパターンについて学びます。

よくある質問

「リアルタイムイベント処理向けKinesis Data Streams」レッスンは無料ですか?

はい。「リアルタイムイベント処理向けKinesis Data Streams」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Cloud & IT Cert Prepコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Cloud & IT Cert Prepコースには全4レッスンが含まれています。

「リアルタイムイベント処理向けKinesis Data Streams」で何を学びますか?

Kinesis Data Streamsで高スループットのイベントストリームを生成・消費し、スループットに応じてシャードを管理して、Lambdaをコンシューマーとして使用します。 ブラウザで直接実行するハンズオンコードでCloud & IT Cert Prepを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

Cloud & IT Cert Prepを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのCloud & IT Cert Prepは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。

「リアルタイムイベント処理向けKinesis Data Streams」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このCloud & IT Cert Prepレッスンでコードを書いて実行できますか?

はい。すべてのCloud & IT Cert Prepレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. EventBridge:イベントバスとルール
  2. Step Functions:サーバーレスワークフローのオーケストレーション
  3. リアルタイムイベント処理向けKinesis Data Streams
  4. コレオグラフィーとオーケストレーションのパターン
← Cloud & IT Cert Prepに戻る