Partitioning Kinesis Records with AWS IoT Payload Values
Use a customer ID from an IoT topic payload as the Kinesis partition key to preserve record order for each customer.
An AWS IoT topic rule can use a value from the message payload as the partition key for Kinesis records. This is useful when downstream processing must preserve the order of records for each customer. The following example uses customer_id as the partition key.
Prerequisites
Install the following tools on your machine:
- AWS SAM CLI
- Python 3.x
Building
The AWS IoT topic rule specifies PartitionKey: ${customer_id} (line 37), which determines how records are distributed across shards.
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: 1The Lambda function that processes the streamed records is a short Python script:
import json
def lambda_handler(event, context): for record in event['Records']: print(json.dumps(record['kinesis']))Build and deploy the stack with SAM:
sam buildsam deploy \ --stack-name aws-iot-kinesis-partition-lambda \ --capabilities CAPABILITY_NAMED_IAMTesting
Test payloads can be sent to the AWS IoT topic directly:
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-outCheck the Lambda logs to confirm that the partitionKey matches the customer_id in your payload.

Cleaning Up
Remove the resources provisioned by this example with:
sam delete --stack-name aws-iot-kinesis-partition-lambdaConclusion
Setting PartitionKey: ${customer_id} in the IoT topic rule sends records for the same customer to the same Kinesis shard without requiring partitioning logic in the Lambda consumer. Because Kinesis uses the partition key to select a shard, records with the same key can be processed in order by downstream consumers.
Load is distributed evenly only when traffic is reasonably balanced across customer_id values. A single high-volume customer can still concentrate traffic on one shard even when the partition key is configured correctly.
Related posts
Streaming OPC UA Data to Kinesis via SiteWise Edge Gateway
Bridging OPC UA telemetry to Kinesis Data Streams through SiteWise Edge Gateway and a custom Greengrass component.
Configuring Record Separators for Kinesis Firehose in AWS IoT Core
Configure a record separator in an IoT Core topic rule's Firehose action so that records stored in S3 are separated by newlines.
Building and Deploying Greengrass Components in a Dockerized Environment
Develop AWS IoT Greengrass components locally using the Greengrass Core Docker image.
Streamlining NDJSON Processing with Kinesis Firehose Dynamic Partitioning
Use Kinesis Data Firehose dynamic partitioning to write NDJSON (Newline Delimited JSON) to S3 without a transformation Lambda function.
Querying S3 Logs with Athena and Kinesis Data Firehose
Deliver logs to S3 with Kinesis Data Firehose and query them in Athena, using custom prefixes to register partitions automatically.
