AWS IoT Core における Kinesis Firehose のレコードセパレーター設定
IoT Core のトピックルールにある Firehose アクションでレコードセパレーターを設定し、S3 に改行区切りでレコードを格納する方法。
IoT Core のトピックルールでは、個々の Firehose レコードの間にレコードセパレーターを挿入できます。これは、後続のコンシューマーが一続きのブロブではなく改行区切りの出力を期待している場合に重要になります。
構築
ポイントとなるのは、13〜14 行目にある Separator: |+ <NEW_LINE> の設定です。
AWSTemplateFormatVersion: "2010-09-09"
Resources: TopicRule: Type: AWS::IoT::TopicRule Properties: RuleName: topic_rule_firehose_separator_test TopicRulePayload: Actions: - Firehose: DeliveryStreamName: !Ref Firehose RoleArn: !GetAtt IamTopicRule.Arn Separator: |+
AwsIotSqlVersion: 2016-03-23 RuleDisabled: false Sql: !Sub SELECT * FROM 'topic_rule_firehose_separator_test'
Firehose: Type: AWS::KinesisFirehose::DeliveryStream Properties: DeliveryStreamName: topic-rule-firehose-separator-test DeliveryStreamType: DirectPut S3DestinationConfiguration: BucketARN: !GetAtt S3.Arn BufferingHints: IntervalInSeconds: 60 SizeInMBs: 5 CompressionFormat: GZIP ErrorOutputPrefix: "error/!{firehose:error-output-type}/!{timestamp:'year='yyyy'/month='MM'/day='dd'/hour='HH}/" Prefix: "success/!{timestamp:'year='yyyy'/month='MM'/day='dd'/hour='HH}/" RoleARN: !GetAtt IamFirehose.Arn
S3: Type: AWS::S3::Bucket Properties: BucketEncryption: ServerSideEncryptionConfiguration: - ServerSideEncryptionByDefault: SSEAlgorithm: AES256 BucketName: topic-rule-firehose-separator-test PublicAccessBlockConfiguration: BlockPublicAcls: TRUE BlockPublicPolicy: TRUE IgnorePublicAcls: TRUE RestrictPublicBuckets: TRUE
IamTopicRule: Type: AWS::IAM::Role Properties: AssumeRolePolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Principal: Service: iot.amazonaws.com Action: sts:AssumeRole Policies: - PolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Action: firehose:PutRecord Resource: - !GetAtt Firehose.Arn PolicyName: policy RoleName: iam-topic-rule
IamFirehose: Type: AWS::IAM::Role Properties: AssumeRolePolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Principal: Service: firehose.amazonaws.com Action: sts:AssumeRole Condition: StringEquals: sts:ExternalId: !Ref AWS::AccountId Policies: - PolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Action: - glue:GetTable - glue:GetTableVersion - glue:GetTableVersions Resource: "*" - Effect: Allow Action: - s3:AbortMultipartUpload - s3:GetBucketLocation - s3:GetObject - s3:ListBucket - s3:ListBucketMultipartUploads - s3:PutObject Resource: - !GetAtt S3.Arn - Fn::Sub: - ${arn}/* - {arn: !GetAtt S3.Arn} - Effect: Allow Action: - lambda:InvokeFunction - lambda:GetFunctionConfiguration Resource: !Sub arn:aws:lambda:${AWS::Region}:${AWS::AccountId}:function:%FIREHOSE_DEFAULT_FUNCTION%:%FIREHOSE_DEFAULT_VERSION% - Effect: Allow Action: - logs:PutLogEvents Resource: - !Sub arn:aws:logs:${AWS::Region}:${AWS::AccountId}:log-group:/aws/kinesisfirehose/topic-rule-firehose-separator-test - Effect: Allow Action: - kinesis:DescribeStream - kinesis:GetShardIterator - kinesis:GetRecords Resource: !Sub arn:aws:kinesis:${AWS::Region}:${AWS::AccountId}:stream/%FIREHOSE_STREAM_NAME% - Effect: Allow Action: - kms:Decrypt Resource: - !Sub arn:aws:kms:${AWS::Region}:${AWS::AccountId}:key/%SSE_KEY_ID% Condition: StringEquals: kms:ViaService: kinesis.%REGION_NAME%.amazonaws.com StringLike: kms:EncryptionContext:aws:kinesis:arn: !Sub arn:aws:kinesis:%REGION_NAME%:${AWS::AccountId}:stream/%FIREHOSE_STREAM_NAME% PolicyName: policy RoleName: iam-firehoseCloudFormation スタックをデプロイします。
aws cloudformation deploy \ --template template.yaml \ --stack-name topic-rule-firehose-separator-test \ --capabilities CAPABILITY_NAMED_IAMトピックルールの設定を確認します。
aws iot get-topic-rule \ --rule-name topic_rule_firehose_separator_test出力の 12 行目に separator の設定が表示されるはずです。
{ "ruleArn": "arn:aws:iot:<YOUR_REGION>:<YOUR_ACCOUNT_ID>:rule/topic_rule_firehose_separator_test", "rule": { "ruleName": "topic_rule_firehose_separator_test", "sql": "SELECT * FROM 'topic_rule_firehose_separator_test'", "createdAt": "2020-05-13T10:29:18+09:00", "actions": [ { "firehose": { "roleArn": "arn:aws:iam::<YOUR_ACCOUNT_ID>:role/iam-topic-rule", "deliveryStreamName": "topic-rule-firehose-separator-test", "separator": "\n" } } ], "ruleDisabled": false, "awsIotSqlVersion": "2016-03-23" }}テスト
トピック topic_rule_firehose_separator_test にテストメッセージを送信します。
aws iot-data publish \ --topic topic_rule_firehose_separator_test \ --payload '{"id": 1, "message": "Hello from AWS IoT"}' \ --cli-binary-format raw-in-base64-out
aws iot-data publish \ --topic topic_rule_firehose_separator_test \ --payload '{"id": 2, "message": "Hello from AWS IoT"}' \ --cli-binary-format raw-in-base64-out生成されたオブジェクトを S3 から取得して確認できます。
# Check an object.$ aws s3 ls topic-rule-firehose-separator-test --recursive2020-05-14 11:30:59 68 success/year=2020/month=05/day=14/hour=02/topic-rule-firehose-separator-test-1-2020-05-14-02-29-57-593d65e5-beb6-47b1-8266-83e869b0cccb.gz
# Download the object.$ aws s3 cp s3://topic-rule-firehose-separator-test/success/year=2020/month=05/day=14/hour=02/topic-rule-firehose-separator-test-1-2020-05-14-02-29-57-593d65e5-beb6-47b1-8266-83e869b0cccb.gz ./result.gz出力には、設定したセパレーターで区切られた 2 件のレコードが含まれているはずです。
# Check the JSON.$ gunzip -c result.gz > result.json$ cat result.json{"id": 1, "message": "Hello from AWS IoT"}{"id": 2, "message": "Hello from AWS IoT"}
# Delete the results.$ rm result.gz result.json$ aws s3 rm s3://topic-rule-firehose-separator-test/success/year=2020/month=05/day=14/hour=02/topic-rule-firehose-separator-test-1-2020-05-14-02-29-57-593d65e5-beb6-47b1-8266-83e869b0cccb.gz後片付け
作業が終わったらスタックを削除します。
aws cloudformation delete-stack --stack-name topic-rule-firehose-separator-testまとめ
IoT トピックルールの Firehose アクションに Separator を追加してデプロイし、ダウンロードした S3 オブジェクトを確認したところ、2 件のテストレコードが一続きのブロブではなく改行区切りの JSON として格納されていました。Separator: |+ は YAML のブロックスカラーで、リテラルな改行に解決されます。このたった一つのプロパティが、きれいに区切られたレコードになるか、後段で独自の分割ロジックが必要な出力になるかの分かれ目です。実際のトラフィックがパイプラインを流れる前に、aws iot get-topic-rule で設定を確認し、さらにダウンロードして gunzip したオブジェクトで実際のセパレーターを確認しておく価値があります。セパレーターが設定されていなくてもエラーにはならず、単に全レコードが黙って連結されてしまうためです。
Related posts
SiteWise Edge GatewayでOPC UAデータをKinesisにストリーミングする
SiteWise Edge Gatewayとカスタムのgreengrassコンポーネントを使って、OPC UAテレメトリをKinesis Data Streamsへ橋渡しします。
Kinesis Firehoseの動的パーティショニングでNDJSON処理を簡素化する
Kinesis Data Firehoseが動的パーティショニングをサポートしたことで、NDJSON(改行区切りJSON)への変換のためだけにLambda関数を用意する必要がなくなりました。
AWS IoT トピックペイロードによる Kinesis シャードの分離
IoT トピックのペイロードを顧客 ID ごとに異なる Kinesis シャードへルーティングし、後続処理でも顧客ごとの順序を保つ方法。
AthenaとKinesis Data FirehoseでS3のログをクエリする
Kinesis Data FirehoseでログをS3に流し込み、Athenaでクエリする。カスタムプレフィックスを使ってパーティションを自動登録させる方法。
Docker環境でGreengrassコンポーネントをビルド・デプロイする
Greengrass Core Dockerイメージを使って、AWS IoT Greengrassコンポーネントをローカルで開発します。
