Kinesis Firehose の動的パーティショニングで NDJSON 処理を簡素化する
Kinesis Data Firehose の動的パーティショニングを使い、変換用の Lambda 関数を用意せずに NDJSON(改行区切り JSON)を S3 へ出力します。
Kinesis Data Firehose は 2021 年に動的パーティショニングをサポートしました。パーティションキーの抽出とレコード区切り文字の追加ができるため、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 を追加すると、整形された NDJSON を出力できます。
MetadataExtractionQuery がレコード構造と一致しない場合、Firehose はレコードをエラー用プレフィックスへ出力します。実際のトラフィックを初めて流すときは、S3 の 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 アプリケーションに組み込みます。
