実務Tips
AWS CLI コマンド
Section titled “AWS CLI コマンド”Kinesis Data Streams
Section titled “Kinesis Data Streams”# ストリーム作成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.jsoncat > 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-deletionKinesis Data Firehose
Section titled “Kinesis Data Firehose”# 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.jsoncat > 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-firehoseKinesis Data Analytics
Section titled “Kinesis Data Analytics”# アプリケーション作成(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"ベストプラクティス
Section titled “ベストプラクティス”シャード設計
Section titled “シャード設計”計算: 必要シャード数 = MAX(書き込みスループット/1MB, 書き込みレコード数/1000, 読み込みスループット/2MB)
例: 書き込み: 5MB/秒 読み込み: 10MB/秒 必要シャード数: MAX(5/1, 10/2) = 5 シャード
パーティションキー: 良い設計: - 高カーディナリティ(ユニーク値が多い) - 均等分散 - 例: user_id, device_id
悪い設計: - 低カーディナリティ - 偏った分散 - 例: 固定値, true/falseプロデューサー最適化
Section titled “プロデューサー最適化”KPL(Kinesis Producer Library): メリット: - 自動バッチ処理 - リトライ機能 - メトリクス収集
設定: Aggregation: - 複数レコードを1つに集約 - スループット向上
Collection: - バッファリング - 遅延と効率のトレードオフ
単純なProducer: - PutRecord / PutRecords API - 直接制御 - Lambda / 小規模アプリケーションコンシューマー最適化
Section titled “コンシューマー最適化”Lambda: 設定: BatchSize: 100-10000 - 大きいほど効率的 - 処理時間考慮
ParallelizationFactor: 1-10 - シャード内並列処理 - レイテンシ削減
MaximumBatchingWindowInSeconds: 0-300 - バッファ時間 - バッチサイズ到達待ち
エラーハンドリング: - BisectBatchOnFunctionError: true - MaximumRetryAttempts: 3 - OnFailureDestination: SQS/SNS
Enhanced Fan-Out: - 専用スループット - 低レイテンシ(70ms) - コンシューマー毎に課金 - 複数コンシューマー推奨Firehose最適化
Section titled “Firehose最適化”バッファ設定: 小データ量: Size: 1MB Interval: 60秒 → レイテンシ優先
大データ量: Size: 128MB Interval: 900秒 → コスト優先
データ変換: - Lambda簡潔に - タイムアウト考慮 - エラーハンドリング
圧縮: - GZIP推奨 - ストレージコスト削減 - Athena互換モニタリング
Section titled “モニタリング”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セキュリティ
Section titled “セキュリティ”暗号化: 転送中: - TLS/SSL必須
保管時: - KMS暗号化有効化 - カスタマー管理キー推奨
アクセス制御: Producer: - kinesis:PutRecord - kinesis:PutRecords
Consumer: - kinesis:GetRecords - kinesis:GetShardIterator - kinesis:DescribeStream
VPCエンドポイント: - プライベート接続 - インターネット不要 - セキュリティ向上コスト最適化
Section titled “コスト最適化”Data Streams: オンデマンド vs プロビジョンド: オンデマンド: - 予測不可能なワークロード - 管理不要 - 従量課金
プロビジョンド: - 予測可能 - コスト効率(高スループット時) - 時間課金
データ保持: - 必要最小限 - デフォルト24時間で十分か確認
Firehose: - バッファサイズ最適化 - Lambda変換効率化 - 圧縮有効化トラブルシューティング
Section titled “トラブルシューティング”よくある問題:
1. ProvisionedThroughputExceeded: 原因: - シャード不足 - パーティションキー偏り
対処: - シャード追加 - パーティションキー見直し - リトライロジック実装
2. Iterator Age増加: 原因: - コンシューマー処理遅延 - Lambda timeout
対処: - Lambda並列度上げ - 処理最適化 - シャード追加
3. Firehose配信失敗: 原因: - S3権限不足 - Lambda変換エラー - バッファタイムアウト
対処: - IAMロール確認 - Lambda CloudWatch Logs確認 - エラーログ確認
4. データ欠損: 原因: - Producer エラー - Consumer エラー - ネットワーク問題
対処: - リトライ実装 - DLQ設定 - CloudTrail確認データリプレイ
Section titled “データリプレイ”用途: - バグ修正後の再処理 - 新しい分析ロジック適用 - 監査・調査
実装: TRIM_HORIZON: - ストリームの最古データから
AT_TIMESTAMP: - 指定時刻から
手順: 1. 保持期間内データ確認 2. シャードイテレータ取得 3. GetRecords で順次取得 4. 再処理実務での重要ポイント:
- シャード設計でスループット確保
- パーティションキーで均等分散
- KPL でプロデューサー最適化
- Lambda バッチサイズ調整
- Enhanced Fan-Out で低レイテンシ
- Firehose バッファ設定最適化
- CloudWatch で継続的監視
- VPCエンドポイントでセキュリティ強化