Kinesis Streams、Firehose、リアルタイム分析
Kinesis Data Streamsでストリーミングデータを取り込み、FirehoseでS3またはRedshiftに配信し、Managed Service for Apache Flinkでリアルタイムに分析します。
「Kinesis Streams、Firehose、リアルタイム分析」はCoddyKit上の無料Cloud & IT Cert Prepレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはCloud & IT Cert Prep学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Cloud & IT Cert Prepコースには全4レッスンが含まれています。
Kinesis ファミリーの概要
Amazon Kinesis は、リアルタイムのストリーミングデータを収集、処理、分析するためのサービス群です。中核となる 3 つのサービスは、Kinesis Data Streams(低レイテンシーのカスタム処理)、Kinesis Data Firehose(S3、Redshift、OpenSearch へのフルマネージド配信)、そしてリアルタイム SQL と Flink 処理に対応するManaged Service for Apache Flink(旧 Kinesis Data Analytics)です。それぞれのサービスは、ストリーミングパイプライン内の異なる用途を対象としています。
Kinesis Data Streams のアーキテクチャ
Kinesis Data Stream は、シャードに分割された、耐久性のある順序付きログです。各シャードは、書き込みスループット 1 MB/秒、読み取りスループット 2 MB/秒を提供します。データレコードはデフォルトで 24 時間保持されます(7 日または 365 日まで延長可能です)。プロデューサーはストリームにレコードを書き込み、コンシューマー(Lambda、KCL アプリケーション、Firehose、Flink など)は 1 つ以上のシャードから並列に読み取ります。レコードは一度書き込まれると変更できません。
# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
--stream-name clickstream \
--shard-count 4
# Put a record into the stream
aws kinesis put-record \
--stream-name clickstream \
--partition-key 'user-123' \
--data 'eyJldmVudCI6ICJjbGljayJ9'シャード、スループット、スケーリング
シャード数によって、ストリーム全体のスループットが決まります。シャードを分割するとスループットを 2 倍にでき、2 つのシャードを統合するとコストを削減できます。Enhanced Fan-Out を使用すると、登録された各コンシューマーに他のコンシューマーとは独立した 2 MB/秒の読み取りスループットを割り当てられます。これにより、複数のアプリケーションが同じストリームを読み取る場合でも、読み取りスロットリングを防げます。GetRecords.IteratorAgeMilliseconds を監視して、コンシューマーの遅延を検出します。
# Split shard to increase throughput
aws kinesis split-shard \
--stream-name clickstream \
--shard-to-split shardId-000000000001 \
--new-starting-hash-key 170141183460469231731687303715884105728
# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
--consumer-name analytics-appKinesis Data Firehose:マネージド配信
Kinesis Data Firehose は、ストリーミングデータを取得、変換し、S3、Amazon Redshift、Amazon OpenSearch Service、Splunk、HTTP エンドポイントなどの送信先に配信するフルマネージドサービスです。管理するシャードはなく、Firehose が自動的にスケールします。バッファサイズ(1~128 MB)とバッファ間隔(60~900 秒)を設定し、どちらかの上限に先に達した時点で Firehose がデータを配信します。
# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
--delivery-stream-name clickstream-to-s3 \
--s3-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"BucketARN": "arn:aws:s3:::my-data-lake-123",
"Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
"BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
"CompressionFormat": "GZIP"
}'Lambda による Firehose データ変換
Firehose は、配信前にレコードの各バッチで Lambda 関数を呼び出し、処理中のデータを変換、拡充、フィルタリングできます。一般的な用途には、JSON から Parquet への変換(Glue スキーマを使用)、PII フィールドのマスキング、価値の低いイベントの削除などがあります。変換に失敗したレコードは、再処理用に別の S3 エラープレフィックスへ任意で送信できるため、データが失われることはありません。
# Lambda transform function signature for Firehose
def lambda_handler(event, context):
output = []
for record in event['records']:
import base64, json
payload = json.loads(base64.b64decode(record['data']))
# Drop events with no user_id
if not payload.get('user_id'):
output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
else:
output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
return {'records': output}Managed Service for Apache Flink
Amazon Managed Service for Apache Flink(旧 Kinesis Data Analytics)は、フルマネージドインフラストラクチャ上で Apache Flink アプリケーションを実行します。スライディングウィンドウ集計、異常検出、イベントシーケンスのパターンマッチング、ストリーミングデータと参照テーブルの結合など、ステートフルなリアルタイム分析に使用します。Flink のコードは Java、Python、Scala で記述でき、Flink がチェックポイントと exactly-once セマンティクスの状態管理を処理します。
# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
url STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);
SELECT
url,
COUNT(*) AS views,
TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);Streams と Firehose の選択
SAA-C03 試験では、判断基準を理解しておきましょう。サブ秒のレイテンシー、複数のコンシューマーによる同時読み取り、または保持期間とリプレイを完全に制御できるカスタム処理ロジックが必要な場合は、Kinesis Data Streams を使用します。最小限のコードでストリーミングデータを S3、Redshift、OpenSearch に確実に配信し、任意の処理中変換と高いレイテンシー(60 秒以上)での自動スケーリングを利用したい場合は、Firehose を使用します。
Kinesis と SQS:試験で頻出する選択
SAA-C03 では、Kinesis と SQS のどちらを選ぶかを問う問題がよく出題されます。主な違いは次のとおりです。Kinesis はシャード内のメッセージの順序を保持し、複数のコンシューマーが同じデータを同時に読み取ることをサポートし、リプレイ用にレコードを保持します。SQS はコンシューマーによる処理後にメッセージを削除するため(リプレイ不可)、FIFO モードで厳密な順序を保証し、マイクロサービスの疎結合に適しています。シナリオにリアルタイム分析やリプレイが含まれている場合は、Kinesis を選択します。
Kinesis プロデューサー:SDK と KPL
Kinesis Producer Library(KPL) は、アプリケーションから Kinesis Data Streams に書き込むための高スループットクライアントです。KPL は複数の小さなレコードを自動的に集約して 1 回の API 呼び出し(最大 1 MB)にまとめ、バックオフを伴う再試行も処理します。これにより、レコードごとの PUT コストを大幅に削減し、シャードあたりのスループットを向上できます。Web クリックストリーム、IoT センサー、ログパイプラインなど、大量のデータを生成するプロデューサーには KPL を使用します。
# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });Firehose から Redshift へのパターン
一般的なアーキテクチャでは、Firehose をマネージドパイプラインとして使用し、Kinesis ストリームまたは直接接続されたプロデューサーからAmazon Redshift にデータを取り込みます。Firehose はまず中間の S3 ステージングバケットにデータを書き込み、次に COPY コマンドを発行してデータを Redshift にロードします。これはストリーミングデータを Redshift に一括ロードする最も効率的な方法です。Redshift に行単位で直接 INSERT すると、行ごとのオーバーヘッドによって非常に遅くなります。
# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
"CopyCommand": {
"DataTableName": "clickevents",
"CopyOptions": "JSON 'auto'"
},
"Username": "firehose_user",
"Password": "{{resolve:secretsmanager:redshift-pw}}",
"S3Configuration": {
"RoleARN": "...",
"BucketARN": "arn:aws:s3:::firehose-staging"
}
}'CloudWatch による Kinesis の監視
Kinesis で重要な CloudWatch メトリクスは次のとおりです。IncomingBytes と IncomingRecords はプロデューサーのスループットを測定し、GetRecords.IteratorAgeMilliseconds はコンシューマーの遅延を測定します(値が大きい場合、コンシューマーが処理に追いついていないことを意味します)。WriteProvisionedThroughputExceeded は、プロデューサーがシャードの上限に達していることを検出します。イテレーターの経過時間とプロビジョンドスループット超過に対して CloudWatch アラームを設定し、Application Auto Scaling によるシャードの Auto Scaling を起動します。
# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
--alarm-name kinesis-high-lag \
--metric-name GetRecords.IteratorAgeMilliseconds \
--namespace AWS/Kinesis \
--dimensions Name=StreamName,Value=clickstream \
--statistic Maximum \
--period 60 \
--threshold 60000 \
--comparison-operator GreaterThanThreshold \
--evaluation-periods 3 \
--alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alertsクイックチェック
このレッスンで学んだ AWS Solutions Architect(SAA-C03)の概念について理解度を確認します。
レッスンのまとめ
このレッスンでは、Kinesis Data Streams が、複数のコンシューマーとリプレイ機能を備えた、耐久性のある順序付きシャード分割ストリーミングを提供すること、Kinesis Data Firehose が、S3、Redshift、OpenSearch へのフルマネージドでコード不要の配信を提供すること、そしてManaged Service for Apache Flink が、ストリーム上でステートフルなリアルタイム分析を可能にすることを学びました。次は、クラウド移行戦略における 7 つの R について学びます。
よくある質問
「Kinesis Streams、Firehose、リアルタイム分析」レッスンは無料ですか?
はい。「Kinesis Streams、Firehose、リアルタイム分析」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Cloud & IT Cert Prepコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Cloud & IT Cert Prepコースには全4レッスンが含まれています。
「Kinesis Streams、Firehose、リアルタイム分析」で何を学びますか?
Kinesis Data Streamsでストリーミングデータを取り込み、FirehoseでS3またはRedshiftに配信し、Managed Service for Apache Flinkでリアルタイムに分析します。 ブラウザで直接実行するハンズオンコードでCloud & IT Cert Prepを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Cloud & IT Cert Prepを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのCloud & IT Cert Prepは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。
「Kinesis Streams、Firehose、リアルタイム分析」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このCloud & IT Cert Prepレッスンでコードを書いて実行できますか?
はい。すべてのCloud & IT Cert Prepレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- S3でデータレイクを構築する
- AWS Glue:ETLとデータカタログ
- Amazon Athena:S3上のサーバーレスSQL
- Kinesis Streams、Firehose、リアルタイム分析