使用AWS Pipes替代Lambda实现DynamoDB流到EventBridge的批处理
使用AWS Pipes替代Lambda实现DynamoDB流到EventBridge的自定义批处理
原方案概述
当前通过DynamoDB流触发Lambda函数,将事件按pk、sk前缀、eventName及字段差异键(diffKeys)分组后推送到EventBridge。原CloudFormation中表和事件源映射配置如下:
"MyTable": { "Type": "AWS::DynamoDB::Table", "Properties": { "AttributeDefinitions": [ { "AttributeName": "pk", "AttributeType": "S" }, { "AttributeName": "sk", "AttributeType": "S" } ], "BillingMode": "PAY_PER_REQUEST", "KeySchema": [ { "AttributeName": "pk", "KeyType": "HASH" }, { "AttributeName": "sk", "KeyType": "RANGE" } ], "GlobalSecondaryIndexes": [], "StreamSpecification": { "StreamViewType": "NEW_AND_OLD_IMAGES" } } }, "MyTableMapping": { "Type": "AWS::Lambda::EventSourceMapping", "Properties": { "FunctionName": { "Ref": "MyTableStreamingFunction" }, "StartingPosition": "LATEST", "MaximumBatchingWindowInSeconds": 1, "EventSourceArn": { "Fn::GetAtt": [ "MyTable", "StreamArn" ] }, "MaximumRetryAttempts": 3 } }
Lambda函数核心逻辑是对DynamoDB流事件进行分组(按pk、sk前缀、eventName、diffKeys),再批量推送到EventBridge。
AWS Pipes替代方案
AWS Pipes可以直接连接DynamoDB流到EventBridge,无需自定义Lambda,通过配置批处理规则和转换模板实现与原Lambda一致的逻辑。
核心配置思路
- 源配置:以DynamoDB流为Pipe源,配置与原EventSourceMapping一致的批处理窗口、重试策略。
- 批处理分组:利用Pipes的
BatchKey特性,将pk、sk前缀、eventName、diffKeys组合为唯一键,实现同组事件合并。 - 转换逻辑:通过Velocity模板计算
diffKeys并构建符合EventBridge要求的事件格式。 - 目标配置:将转换后的事件推送到EventBridge。
具体CloudFormation配置
"MyTable": { "Type": "AWS::DynamoDB::Table", "Properties": { // 保留原表配置,StreamSpecification必须保留 "AttributeDefinitions": [ { "AttributeName": "pk", "AttributeType": "S" }, { "AttributeName": "sk", "AttributeType": "S" } ], "BillingMode": "PAY_PER_REQUEST", "KeySchema": [ { "AttributeName": "pk", "KeyType": "HASH" }, { "AttributeName": "sk", "KeyType": "RANGE" } ], "GlobalSecondaryIndexes": [], "StreamSpecification": { "StreamViewType": "NEW_AND_OLD_IMAGES" } } }, "MyDynamoToEventBridgePipe": { "Type": "AWS::Pipes::Pipe", "Properties": { "Name": "DynamoDBToEventBridgePipe", "Description": "Pipe from DynamoDB Stream to EventBridge with custom batching", "Source": { "Fn::GetAtt": [ "MyTable", "StreamArn" ] }, "SourceParameters": { "DynamoDBStreamParameters": { "StartingPosition": "LATEST", "BatchSize": 100, "MaximumBatchingWindowInSeconds": 1, "MaximumRetryAttempts": 3, "ParallelizationFactor": 1 }, "BatchParameters": { "BatchSize": 100, "BatchKey": "concat(dynamodb.Keys.pk.S, '/', split(dynamodb.Keys.sk.S, '#')[0], '/', eventName, '/', join('|', sort(diffKeys)))", "MaximumBatchingWindowInSeconds": 1 } }, "Transformer": { "InputTemplate": "{\n \"Source\": \"MyDynamoToEventBridgePipe\",\n \"DetailType\": \"<eventName>\",\n \"Detail\": {\n \"pk\": \"<dynamodb.Keys.pk.S>\",\n \"sk\": \"<split(dynamodb.Keys.sk.S, '#')[0]>\",\n \"eventName\": \"<eventName>\",\n \"diffKeys\": <#if dynamodb.NewImage && dynamodb.OldImage>\n <#assign newImage = {}>\n <#list dynamodb.NewImage as k, v>\n <#assign newImage[k] = v?values[0]>\n </#list>\n <#assign oldImage = {}>\n <#list dynamodb.OldImage as k, v>\n <#assign oldImage[k] = v?values[0]>\n </#list>\n <#assign diffKeys = []>\n <#list newImage as k, v>\n <#if !(k in oldImage) || newImage[k] != oldImage[k]>\n <#assign diffKeys.add(k)>\n </#if>\n </#list>\n ${diffKeys?sort?json}\n <#else>\n []\n </#if>,\n \"records\": <Records?json>\n }\n}" }, "Target": "arn:aws:events:<region>:<account-id>:event-bus/default", "TargetParameters": { "EventBridgeEventBusParameters": { "DetailType": "<DetailType>", "Source": "<Source>" } }, "RoleArn": { "Fn::GetAtt": [ "PipeExecutionRole", "Arn" ] } } }, "PipeExecutionRole": { "Type": "AWS::IAM::Role", "Properties": { "AssumeRolePolicyDocument": { "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Principal": { "Service": "pipes.amazonaws.com" }, "Action": "sts:AssumeRole" } ] }, "Policies": [ { "PolicyName": "PipeExecutionPolicy", "PolicyDocument": { "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "dynamodb:DescribeStream", "dynamodb:GetRecords", "dynamodb:GetShardIterator", "dynamodb:ListStreams" ], "Resource": { "Fn::GetAtt": [ "MyTable", "StreamArn" ] } }, { "Effect": "Allow", "Action": [ "events:PutEvents" ], "Resource": "arn:aws:events:<region>:<account-id>:event-bus/default" } ] } } ] } }
关键逻辑对应说明
- 批处理分组:
原Lambda通过Key类组合pk、sk前缀、eventName、diffKeys作为分组键,Pipes中通过BatchKey使用表达式生成相同的唯一键,实现同组事件合并。 - diffKeys计算:
原Lambda的diff_keys函数对比NewImage和OldImage的差异字段,Pipes的InputTemplate使用Velocity模板遍历新旧图像,计算差异并排序,与原逻辑完全一致。 - 事件格式转换:
原Lambda的Entry.entry构建EventBridge事件结构,Pipes的InputTemplate直接生成符合要求的JSON格式,包含Source、DetailType、Detail等字段。 - 批处理窗口与重试:
Pipes的MaximumBatchingWindowInSeconds和MaximumRetryAttempts与原Lambda的EventSourceMapping配置一致,保证相同的事件处理延迟和重试策略。
注意事项
- 需替换配置中的
<region>、<account-id>为实际AWS区域和账号ID。 BatchSize可根据原Lambda的BATCH_SIZE环境变量调整,最大支持1000条事件。
内容的提问来源于stack exchange,提问作者Justin
相关产品推荐
相关产品推荐

