AWS IoT のペイロード値で Kinesis レコードをパーティショニングする
IoT トピックのペイロードに含まれる顧客 ID を Kinesis のパーティションキーとして使い、顧客ごとのレコード順序を維持する方法を解説します。
AWS IoT のトピックルールでは、メッセージのペイロードに含まれる値を Kinesis レコードのパーティションキーとして使用できます。これは、後続処理で顧客ごとのレコード順序を維持したい場合に役立ちます。以下の例では、customer_id をパーティションキーとして使用します。
前提条件
以下のツールをあらかじめインストールします。
- AWS SAM CLI
- Python 3.x
構築
AWS IoT のトピックルールでは PartitionKey: ${customer_id}(37 行目)を指定しており、これによってレコードがどのようにシャードへ振り分けられるかが決まります。
AWSTemplateFormatVersion: '2010-09-09'Transform: AWS::Serverless-2016-10-31
Resources: Lambda: Type: AWS::Serverless::Function Properties: CodeUri: ./src/ Events: Kinesis: Type: Kinesis Properties: BatchSize: 100 BisectBatchOnFunctionError: true Enabled: true StartingPosition: LATEST Stream: !GetAtt Kinesis.Arn FunctionName: aws_iot_kinesis_partition_lambda Handler: lambda_function.lambda_handler Role: !GetAtt IamRole.Arn Runtime: python3.8
Kinesis: Type: AWS::Kinesis::Stream Properties: Name: aws_iot_kinesis_partition_stream RetentionPeriodHours: 24 ShardCount: 1
TopicRule: Type: AWS::IoT::TopicRule Properties: RuleName: aws_iot_kinesis_partition_topic_rule TopicRulePayload: Actions: - Kinesis: PartitionKey: ${customer_id} RoleArn: !GetAtt IamRole.Arn StreamName: aws_iot_kinesis_partition_stream AwsIotSqlVersion: 2016-03-23 RuleDisabled: false Sql: SELECT * FROM 'aws-iot-kinesis-partition-topic'
IamRole: Type: AWS::IAM::Role Properties: AssumeRolePolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Principal: Service: lambda.amazonaws.com Action: sts:AssumeRole - Effect: Allow Principal: Service: iot.amazonaws.com Action: sts:AssumeRole ManagedPolicyArns: - arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole Policies: - PolicyDocument: Version: '2012-10-17' Statement: - Effect: Allow Action: - kinesis:GetShardIterator - kinesis:GetRecords - kinesis:DescribeStream - kinesis:PutRecord Resource: - !GetAtt Kinesis.Arn PolicyName: policy
LogGroup: Type: AWS::Logs::LogGroup Properties: LogGroupName: /aws/lambda/aws_iot_kinesis_partition_lambda RetentionInDays: 1ストリームされたレコードを処理する Lambda 関数は、短い Python スクリプトです。
import json
def lambda_handler(event, context): for record in event['Records']: print(json.dumps(record['kinesis']))SAM を使ってスタックをビルド・デプロイします。
sam buildsam deploy \ --stack-name aws-iot-kinesis-partition-lambda \ --capabilities CAPABILITY_NAMED_IAMテスト
テスト用のペイロードは、AWS IoT のトピックに直接送信できます。
aws iot-data publish \ --topic aws-iot-kinesis-partition-topic \ --payload '{"customer_id": 1, "message": "Hello from AWS IoT"}' \ --cli-binary-format raw-in-base64-outLambda のログを確認し、partitionKey がペイロード内の customer_id と一致していることを確認します。

後片付け
この例で作成したリソースは、以下のコマンドで削除します。
sam delete --stack-name aws-iot-kinesis-partition-lambdaまとめ
IoT トピックルールに PartitionKey: ${customer_id} を設定すると、Lambda コンシューマーにパーティショニングロジックを実装せずに、同じ顧客のレコードを同じ Kinesis シャードへ送信できます。Kinesis はパーティションキーを基にシャードを選択するため、後続のコンシューマーは同じキーを持つレコードを順番に処理できます。
負荷が均等に分散されるのは、customer_id ごとのトラフィック量がおおむね均等な場合です。特定の顧客のトラフィック量が多いと、パーティションキーを正しく設定していても、1 つのシャードに負荷が集中する可能性があります。
Related posts
SiteWise Edge Gateway で OPC UA データを Kinesis にストリーミングする
SiteWise Edge Gateway とカスタムの Greengrass コンポーネントを使い、OPC UA テレメトリを Kinesis Data Streams へ橋渡しします。
AWS IoT Core における Kinesis Firehose のレコード区切り設定
IoT Core のトピックルールにある Firehose アクションで区切り文字を設定し、S3 にレコードを改行区切りで格納する方法を解説します。
Docker 環境で Greengrass コンポーネントをビルド・デプロイする
Greengrass Core Docker イメージを使い、AWS IoT Greengrass コンポーネントをローカルで開発します。
Kinesis Firehose の動的パーティショニングで NDJSON 処理を簡素化する
Kinesis Data Firehose の動的パーティショニングを使い、変換用の Lambda 関数を用意せずに NDJSON(改行区切り JSON)を S3 へ出力します。
Athena と Kinesis Data Firehose で S3 のログをクエリする
Kinesis Data Firehose でログを S3 に配信し、Athena でクエリします。カスタムプレフィックスを使ってパーティションを自動登録する方法を解説します。
