Kinesis Firehoseの動的パーティショニングでNDJSON処理を簡素化する
Kinesis Data Firehoseが動的パーティショニングをサポートしたことで、NDJSON(改行区切りJSON)への変換のためだけにLambda関数を用意する必要がなくなりました。
Kinesis Data Firehoseは最近、動的パーティショニングのサポートを追加しました。これにより、NDJSON(改行区切りJSON)への変換のためだけにLambda関数を用意する必要がなくなりました。
構築
動的パーティショニングの設定は、以下のCloudFormationテンプレートに記述されています。
AWSTemplateFormatVersion: 2010-09-09Description: Kinesis Data Firehose streaming NDJSON sample with dynamic partitioningResources: KinesisFirehoseDeliveryStream: Type: AWS::KinesisFirehose::DeliveryStream Properties: DeliveryStreamName: ndjson-firehose DeliveryStreamType: DirectPut ExtendedS3DestinationConfiguration: BucketARN: !GetAtt S3Bucket.Arn BufferingHints: IntervalInSeconds: 60 DynamicPartitioningConfiguration: Enabled: true Prefix: "success/user_id=!{partitionKeyFromQuery:user_id}/year=!{timestamp:yyyy}/month=!{timestamp:MM}/day=!{timestamp:dd}/hour=!{timestamp:HH}/" ErrorOutputPrefix: "error/year=!{timestamp:yyyy}/month=!{timestamp:MM}/day=!{timestamp:dd}/hour=!{timestamp:HH}/type=!{firehose:error-output-type}/" ProcessingConfiguration: Enabled: true Processors: - Type: AppendDelimiterToRecord Parameters: - ParameterName: Delimiter ParameterValue: '\\n' - Type: MetadataExtraction Parameters: - ParameterName: MetadataExtractionQuery ParameterValue: '{user_id: .user_id}' - ParameterName: JsonParsingEngine ParameterValue: JQ-1.6 RoleARN: !GetAtt IAMRoleKinesisFirehose.Arn
S3Bucket: Type: AWS::S3::Bucket Properties: BucketName: !Sub ${AWS::Region}-${AWS::AccountId}-ndjson-s3 BucketEncryption: ServerSideEncryptionConfiguration: - ServerSideEncryptionByDefault: SSEAlgorithm: AES256 PublicAccessBlockConfiguration: BlockPublicAcls: true BlockPublicPolicy: true IgnorePublicAcls: true RestrictPublicBuckets: true
IAMRoleKinesisFirehose: Type: AWS::IAM::Role Properties: RoleName: ndjson-firehose-role AssumeRolePolicyDocument: Version: 2012-10-17 Statement: - Effect: Allow Principal: Service: firehose.amazonaws.com Action: sts:AssumeRole MaxSessionDuration: 3600 Policies: - PolicyName: policy1 PolicyDocument: Version: 2012-10-17 Statement: - Effect: Allow Action: - glue:GetTable - glue:GetTableVersion - glue:GetTableVersions Resource: - !Sub arn:aws:glue:${AWS::Region}:${AWS::AccountId}:catalog - !Sub arn:aws:glue:${AWS::Region}:${AWS::AccountId}:database/%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER% - !Sub arn:aws:glue:${AWS::Region}:${AWS::AccountId}:table/%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER%/%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER% - Effect: Allow Action: - s3:AbortMultipartUpload - s3:GetBucketLocation - s3:GetObject - s3:ListBucket - s3:ListBucketMultipartUploads - s3:PutObject Resource: - !GetAtt S3Bucket.Arn - !Sub ${S3Bucket.Arn}/* - Effect: Allow Action: - lambda:InvokeFunction - lambda:GetFunctionConfiguration Resource: !Sub arn:aws:lambda:${AWS::Region}:${AWS::AccountId}:function:%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER% - Effect: Allow Action: - kms:GenerateDataKey - kms:Decrypt Resource: - !Sub arn:aws:kms:${AWS::Region}:${AWS::AccountId}:key/%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER% Condition: StringEquals: kms:ViaService: !Sub s3.${AWS::Region}.amazonaws.com StringLike: kms:EncryptionContext:aws:s3:arn: - arn:aws:s3:::%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER%/* - arn:aws:s3:::%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER% - Effect: Allow Action: - logs:PutLogEvents Resource: - !Sub arn:aws:logs:${AWS::Region}:${AWS::AccountId}:log-group:%FIREHOSE_POLICY_TEMPLATE_PLACEHOLDER%:log-stream:*スタックをデプロイします。
aws cloudformation deploy \ --template-file stack.yaml \ --stack-name firehose-ndjson-sample \ --s3-bucket <YOUR_S3_BUCKET> \ --s3-prefix firehose-ndjson-sample/$(date +%Y/%m/%d/%H) \ --capabilities CAPABILITY_NAMED_IAMテスト
公式ドキュメントに記載の通り、データブロブは送信前にBase64エンコードする必要があります。
データブロブは、シリアライズ時にBase64エンコードされます。Base64エンコード前のデータブロブの最大サイズは1,000 KiBです。
サンプルのJSONを以下のようにエンコードします。
$ echo -n '{"user_id": 1, "message": "Hello"}' | base64eyJ1c2VyX2lkIjogMSwgIm1lc3NhZ2UiOiAiSGVsbG8ifQ==
$ echo -n '{"user_id": 1, "message": "World"}' | base64eyJ1c2VyX2lkIjogMSwgIm1lc3NhZ2UiOiAiV29ybGQifQ==エンコードしたレコードは、AWS CLI経由でFirehoseストリームに送信できます。
echo '{ "DeliveryStreamName": "ndjson-firehose", "Records": [ {"Data": "eyJ1c2VyX2lkIjogMSwgIm1lc3NhZ2UiOiAiSGVsbG8ifQ=="}, {"Data": "eyJ1c2VyX2lkIjogMSwgIm1lc3NhZ2UiOiAiV29ybGQifQ=="} ]}' > input.jsonaws firehose put-record-batch --cli-input-json file://~/input.jsonS3バケット内のオブジェクトを確認します。
$ aws s3 ls --recursive s3://<AWS_REGION>-<AWS_ACCOUNT>-ndjson-s3/success/2022-07-26 23:52:47 69 success/user_id=1/year=2022/month=07/day=26/hour=14/ndjson-firehose-4-2022-07-26-14-50-22-b95618c0-e518-3b66-b06c-693b059cc751
$ aws s3 cp s3://<AWS_REGION>-<AWS_ACCOUNT>-ndjson-s3/success/user_id=1/year=2022/month=07/day=26/hour=14/ndjson-firehose-4-2022-07-26-14-50-22-b95618c0-e518-3b66-b06c-693b059cc751 ./download: s3://<AWS_REGION>-<AWS_ACCOUNT>-ndjson-s3/success/user_id=1/year=2022/month=07/day=26/hour=14/ndjson-firehose-4-2022-07-26-14-50-22-b95618c0-e518-3b66-b06c-693b059cc751 to ./ndjson-firehose-4-2022-07-26-14-50-22-b95618c0-e518-3b66-b06c-693b059cc751
$ cat -n ndjson-firehose-4-2022-07-26-14-50-22-b95618c0-e518-3b66-b06c-693b059cc751 1 {"user_id": 1, "message": "Hello"} 2 {"user_id": 1, "message": "World"}後片付け
この例で作成したリソースは、以下で削除できます。
aws s3 rm --recursive s3://<AWS_REGION>-<AWS_ACCOUNT>-ndjson-s3/aws cloudformation delete-stack --stack-name firehose-ndjson-sampleまとめ
Firehoseの配信ストリームで動的パーティショニングを有効にしたところ、パイプラインに変換用のLambdaを組み込むことなく、パーティション分けされた改行区切りのJSONをS3に出力できました。DynamicPartitioningConfiguration: Enabled: trueとMetadataExtractionプロセッサを組み合わせることで、Prefixテンプレートから!{partitionKeyFromQuery:user_id}を直接参照できるようになります。これは、JQクエリ{user_id: .user_id}によって各レコードのJSON本文からパーティションキーを取り出す仕組みで、事前にLambda関数で計算しておく必要がありません。これにAppendDelimiterToRecordプロセッサの\\n区切り文字を組み合わせることで、これまで変換用Lambdaが必要だったもう半分の役割、つまりきれいなNDJSON出力の生成も担えます。実際のレコード形式と一致しないMetadataExtractionQueryはエラーを発生させず、黙ってエラー用プレフィックスにルーティングされるため、実際のトラフィックを初めて流す際には、success/と併せてerror/パス配下も確認しておく価値があります。
Related posts
AWS IoT Core における Kinesis Firehose のレコードセパレーター設定
IoT Core のトピックルールにある Firehose アクションでレコードセパレーターを設定し、S3 に改行区切りでレコードを格納する方法。
AthenaとKinesis Data FirehoseでS3のログをクエリする
Kinesis Data FirehoseでログをS3に流し込み、Athenaでクエリする。カスタムプレフィックスを使ってパーティションを自動登録させる方法。
SiteWise Edge GatewayでOPC UAデータをKinesisにストリーミングする
SiteWise Edge Gatewayとカスタムのgreengrassコンポーネントを使って、OPC UAテレメトリをKinesis Data Streamsへ橋渡しします。
AWS IoT トピックペイロードによる Kinesis シャードの分離
IoT トピックのペイロードを顧客 ID ごとに異なる Kinesis シャードへルーティングし、後続処理でも顧客ごとの順序を保つ方法。
Cognito User PoolsとOIDCでSlackサインインを実装する
Cognito user poolをOIDC経由でSlackと連携させ、"Sign in with Slack"をAmplifyでNext.jsアプリに組み込みます。
