Skip to content

実務Tips

Terminal window
# ストリーム作成
aws kinesis create-stream \
--stream-name my-stream \
--shard-count 2 \
--region ap-northeast-1
# ストリーム作成(オンデマンドモード)
aws kinesis create-stream \
--stream-name my-stream \
--stream-mode-details StreamMode=ON_DEMAND
# ストリーム一覧
aws kinesis list-streams
# ストリーム詳細
aws kinesis describe-stream \
--stream-name my-stream
# ストリーム詳細(拡張情報)
aws kinesis describe-stream-summary \
--stream-name my-stream
# レコード送信(単一)
aws kinesis put-record \
--stream-name my-stream \
--partition-key user123 \
--data '{"timestamp":"2024-01-01T00:00:00Z","user_id":"user123","event":"login"}' \
--cli-binary-format raw-in-base64-out
# レコード送信(バッチ)
aws kinesis put-records \
--stream-name my-stream \
--records file://records.json
# records.json
cat > records.json << 'EOF'
[
{
"Data": "{\"user_id\":\"user1\",\"event\":\"login\"}",
"PartitionKey": "user1"
},
{
"Data": "{\"user_id\":\"user2\",\"event\":\"logout\"}",
"PartitionKey": "user2"
}
]
EOF
# レコード取得(シャードイテレータ取得)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name my-stream \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' \
--output text)
# レコード取得
aws kinesis get-records \
--shard-iterator $SHARD_ITERATOR
# シャード分割
aws kinesis split-shard \
--stream-name my-stream \
--shard-to-split shardId-000000000000 \
--new-starting-hash-key 170141183460469231731687303715884105728
# シャード統合
aws kinesis merge-shards \
--stream-name my-stream \
--shard-to-merge shardId-000000000000 \
--adjacent-shard-to-merge shardId-000000000001
# データ保持期間変更
aws kinesis increase-stream-retention-period \
--stream-name my-stream \
--retention-period-hours 168 # 7日
aws kinesis decrease-stream-retention-period \
--stream-name my-stream \
--retention-period-hours 24
# Enhanced Fan-Out コンシューマー登録
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:ap-northeast-1:123456789012:stream/my-stream \
--consumer-name my-consumer
# コンシューマー一覧
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:ap-northeast-1:123456789012:stream/my-stream
# ストリーム削除
aws kinesis delete-stream \
--stream-name my-stream \
--enforce-consumer-deletion
Terminal window
# Delivery Stream作成(S3配信)
aws firehose create-delivery-stream \
--delivery-stream-name my-firehose \
--delivery-stream-type DirectPut \
--s3-destination-configuration file://s3-config.json
# s3-config.json
cat > s3-config.json << 'EOF'
{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseDeliveryRole",
"BucketARN": "arn:aws:s3:::my-data-bucket",
"Prefix": "data/year=!{timestamp:yyyy}/month=!{timestamp:MM}/day=!{timestamp:dd}/",
"ErrorOutputPrefix": "errors/year=!{timestamp:yyyy}/month=!{timestamp:MM}/day=!{timestamp:dd}/!{firehose:error-output-type}",
"BufferingHints": {
"SizeInMBs": 5,
"IntervalInSeconds": 300
},
"CompressionFormat": "GZIP",
"EncryptionConfiguration": {
"KMSEncryptionConfig": {
"AWSKMSKeyARN": "arn:aws:kms:ap-northeast-1:123456789012:key/xxxxx"
}
},
"CloudWatchLoggingOptions": {
"Enabled": true,
"LogGroupName": "/aws/kinesisfirehose/my-firehose",
"LogStreamName": "S3Delivery"
}
}
EOF
# Delivery Stream 一覧
aws firehose list-delivery-streams
# Delivery Stream 詳細
aws firehose describe-delivery-stream \
--delivery-stream-name my-firehose
# レコード送信
aws firehose put-record \
--delivery-stream-name my-firehose \
--record '{"Data":"eyJ1c2VyX2lkIjoidXNlcjEyMyIsImV2ZW50IjoibG9naW4ifQ=="}'
# レコード送信(バッチ)
aws firehose put-record-batch \
--delivery-stream-name my-firehose \
--records file://firehose-records.json
# Delivery Stream 更新
aws firehose update-destination \
--delivery-stream-name my-firehose \
--current-delivery-stream-version-id 1 \
--destination-id destinationId-000000000001 \
--s3-destination-update file://s3-update.json
# Delivery Stream 削除
aws firehose delete-delivery-stream \
--delivery-stream-name my-firehose
Terminal window
# アプリケーション作成(SQL)
aws kinesisanalytics create-application \
--application-name my-analytics-app \
--inputs file://input-config.json \
--outputs file://output-config.json \
--application-code file://sql-code.sql
# アプリケーション一覧
aws kinesisanalytics list-applications
# アプリケーション詳細
aws kinesisanalytics describe-application \
--application-name my-analytics-app
# アプリケーション開始
aws kinesisanalytics start-application \
--application-name my-analytics-app \
--input-configurations Id=1.1,InputStartingPositionConfiguration={InputStartingPosition=NOW}
# アプリケーション停止
aws kinesisanalytics stop-application \
--application-name my-analytics-app
# アプリケーション削除
aws kinesisanalytics delete-application \
--application-name my-analytics-app \
--create-timestamp "2024-01-01T00:00:00Z"
計算:
必要シャード数 = MAX(書き込みスループット/1MB, 書き込みレコード数/1000, 読み込みスループット/2MB)
:
書き込み: 5MB/秒
読み込み: 10MB/秒
必要シャード数: MAX(5/1, 10/2) = 5 シャード
パーティションキー:
良い設計:
- 高カーディナリティ(ユニーク値が多い)
- 均等分散
- : user_id, device_id
悪い設計:
- 低カーディナリティ
- 偏った分散
- : 固定値, true/false
KPL(Kinesis Producer Library):
メリット:
- 自動バッチ処理
- リトライ機能
- メトリクス収集
設定:
Aggregation:
- 複数レコードを1つに集約
- スループット向上
Collection:
- バッファリング
- 遅延と効率のトレードオフ
単純なProducer:
- PutRecord / PutRecords API
- 直接制御
- Lambda / 小規模アプリケーション
Lambda:
設定:
BatchSize: 100-10000
- 大きいほど効率的
- 処理時間考慮
ParallelizationFactor: 1-10
- シャード内並列処理
- レイテンシ削減
MaximumBatchingWindowInSeconds: 0-300
- バッファ時間
- バッチサイズ到達待ち
エラーハンドリング:
- BisectBatchOnFunctionError: true
- MaximumRetryAttempts: 3
- OnFailureDestination: SQS/SNS
Enhanced Fan-Out:
- 専用スループット
- 低レイテンシ(70ms)
- コンシューマー毎に課金
- 複数コンシューマー推奨
バッファ設定:
小データ量:
Size: 1MB
Interval: 60秒
→ レイテンシ優先
大データ量:
Size: 128MB
Interval: 900秒
→ コスト優先
データ変換:
- Lambda簡潔に
- タイムアウト考慮
- エラーハンドリング
圧縮:
- GZIP推奨
- ストレージコスト削減
- Athena互換
CloudWatch Metrics:
Data Streams:
重要メトリクス:
- IncomingBytes / IncomingRecords
- GetRecords.IteratorAgeMilliseconds
- WriteProvisionedThroughputExceeded
- ReadProvisionedThroughputExceeded
Firehose:
重要メトリクス:
- IncomingRecords / IncomingBytes
- DeliveryToS3.Success
- DataFreshness(鮮度)
アラーム:
- Iterator Age > 60000ms(遅延)
- ProvisionedThroughputExceeded > 0(スロットリング)
- Lambda Errors > 0
- Firehose DeliveryToS3.Success < 1
暗号化:
転送中:
- TLS/SSL必須
保管時:
- KMS暗号化有効化
- カスタマー管理キー推奨
アクセス制御:
Producer:
- kinesis:PutRecord
- kinesis:PutRecords
Consumer:
- kinesis:GetRecords
- kinesis:GetShardIterator
- kinesis:DescribeStream
VPCエンドポイント:
- プライベート接続
- インターネット不要
- セキュリティ向上
Data Streams:
オンデマンド vs プロビジョンド:
オンデマンド:
- 予測不可能なワークロード
- 管理不要
- 従量課金
プロビジョンド:
- 予測可能
- コスト効率(高スループット時)
- 時間課金
データ保持:
- 必要最小限
- デフォルト24時間で十分か確認
Firehose:
- バッファサイズ最適化
- Lambda変換効率化
- 圧縮有効化
よくある問題:
1. ProvisionedThroughputExceeded:
原因:
- シャード不足
- パーティションキー偏り
対処:
- シャード追加
- パーティションキー見直し
- リトライロジック実装
2. Iterator Age増加:
原因:
- コンシューマー処理遅延
- Lambda timeout
対処:
- Lambda並列度上げ
- 処理最適化
- シャード追加
3. Firehose配信失敗:
原因:
- S3権限不足
- Lambda変換エラー
- バッファタイムアウト
対処:
- IAMロール確認
- Lambda CloudWatch Logs確認
- エラーログ確認
4. データ欠損:
原因:
- Producer エラー
- Consumer エラー
- ネットワーク問題
対処:
- リトライ実装
- DLQ設定
- CloudTrail確認
用途:
- バグ修正後の再処理
- 新しい分析ロジック適用
- 監査・調査
実装:
TRIM_HORIZON:
- ストリームの最古データから
AT_TIMESTAMP:
- 指定時刻から
手順:
1. 保持期間内データ確認
2. シャードイテレータ取得
3. GetRecords で順次取得
4. 再処理

実務での重要ポイント:

  • シャード設計でスループット確保
  • パーティションキーで均等分散
  • KPL でプロデューサー最適化
  • Lambda バッチサイズ調整
  • Enhanced Fan-Out で低レイテンシ
  • Firehose バッファ設定最適化
  • CloudWatch で継続的監視
  • VPCエンドポイントでセキュリティ強化