You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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一致的逻辑。

核心配置思路

  1. 源配置:以DynamoDB流为Pipe源,配置与原EventSourceMapping一致的批处理窗口、重试策略。
  2. 批处理分组:利用Pipes的BatchKey特性,将pk、sk前缀、eventName、diffKeys组合为唯一键,实现同组事件合并。
  3. 转换逻辑:通过Velocity模板计算diffKeys并构建符合EventBridge要求的事件格式。
  4. 目标配置:将转换后的事件推送到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"
            }
          ]
        }
      }
    ]
  }
}

关键逻辑对应说明

  1. 批处理分组:
    原Lambda通过Key类组合pk、sk前缀、eventName、diffKeys作为分组键,Pipes中通过BatchKey使用表达式生成相同的唯一键,实现同组事件合并。
  2. diffKeys计算:
    原Lambda的diff_keys函数对比NewImage和OldImage的差异字段,Pipes的InputTemplate使用Velocity模板遍历新旧图像,计算差异并排序,与原逻辑完全一致。
  3. 事件格式转换:
    原Lambda的Entry.entry构建EventBridge事件结构,Pipes的InputTemplate直接生成符合要求的JSON格式,包含Source、DetailType、Detail等字段。
  4. 批处理窗口与重试:
    Pipes的MaximumBatchingWindowInSeconds和MaximumRetryAttempts与原Lambda的EventSourceMapping配置一致,保证相同的事件处理延迟和重试策略。

注意事项

  • 需替换配置中的<region>、<account-id>为实际AWS区域和账号ID。
  • BatchSize可根据原Lambda的BATCH_SIZE环境变量调整,最大支持1000条事件。

内容的提问来源于stack exchange,提问作者Justin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 14:55:54