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

AWS EventBridge Pipes:能否对事件数组应用Enrichment输入转换器?

如何在EventBridge Pipes中批量解码Kinesis事件数组的data字段?

当前架构:

源过滤增强目标
Kinesis流 →未配置 →Step Functions →SNS Topic

当Kinesis流批处理大小设为大于1时,EventBridge Pipes会将多个Kinesis记录打包成数组传递给Step Functions。单个事件时,输入转换器能正常解码base64格式的data字段,但数组场景下原转换器失效,仅指定第一个元素只能处理单条记录,要实现批量解码每个事件的data字段,可通过以下两种方式实现:

方式1:仅提取所有解码后的data数组

如果只需要每个事件解码后的data内容,使用JSONPath的数组投影语法:

转换器配置:

{
  "decodedData": <$[*].data>
}

输入示例(事件数组):

[
  {
    "kinesisSchemaVersion": "1.0",
    "partitionKey": "1",
    "data": "SGVsbG8sIHRoaXMgaXMgYSB0ZXN0Lg==",
    "approximateArrivalTimestamp": 1545084650.987,
    "eventSource": "aws:kinesis",
    "eventVersion": "1.0",
    "eventID": "shardId-000000000006:49590338271490256608559692538361571095921575989136588898",
    "eventName": "aws:kinesis:record",
    "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
    "awsRegion": "us-east-2",
    "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
  },
  {
    "kinesisSchemaVersion": "1.0",
    "partitionKey": "2",
    "data": "VGhpcyBpcyBhbm90aGVyIHRlc3Qu",
    "approximateArrivalTimestamp": 1545084651.123,
    "eventSource": "aws:kinesis",
    "eventVersion": "1.0",
    "eventID": "shardId-000000000006:49590338271490256608559692538361571095921575989136588900",
    "eventName": "aws:kinesis:record",
    "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
    "awsRegion": "us-east-2",
    "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
  }
]

输出结果:

{
  "decodedData": [
    "Hello, this is a test.",
    "This is another test."
  ]
}

方式2:保留完整事件结构并解码data字段

如果需要保留每个Kinesis事件的完整元数据,仅解码data字段,使用JSONPath的数组投影+对象构造语法:

转换器配置:

{
  "records": <$[*].{
    "kinesisSchemaVersion": kinesisSchemaVersion,
    "partitionKey": partitionKey,
    "sequenceNumber": sequenceNumber,
    "data": data,
    "approximateArrivalTimestamp": approximateArrivalTimestamp,
    "eventSource": eventSource,
    "eventVersion": eventVersion,
    "eventID": eventID,
    "eventName": eventName,
    "invokeIdentityArn": invokeIdentityArn,
    "awsRegion": awsRegion,
    "eventSourceARN": eventSourceARN
  }>
}

输出结果:

{
  "records": [
    {
      "kinesisSchemaVersion": "1.0",
      "partitionKey": "1",
      "sequenceNumber": "49590338271490256608559692538361571095921575989136588898",
      "data": "Hello, this is a test.",
      "approximateArrivalTimestamp": 1545084650.987,
      "eventSource": "aws:kinesis",
      "eventVersion": "1.0",
      "eventID": "shardId-000000000006:49590338271490256608559692538361571095921575989136588898",
      "eventName": "aws:kinesis:record",
      "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
      "awsRegion": "us-east-2",
      "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
    },
    {
      "kinesisSchemaVersion": "1.0",
      "partitionKey": "2",
      "sequenceNumber": "49590338271490256608559692538361571095921575989136588900",
      "data": "This is another test.",
      "approximateArrivalTimestamp": 1545084651.123,
      "eventSource": "aws:kinesis",
      "eventVersion": "1.0",
      "eventID": "shardId-000000000006:49590338271490256608559692538361571095921575989136588900",
      "eventName": "aws:kinesis:record",
      "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
      "awsRegion": "us-east-2",
      "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
    }
  ]
}

关键说明

EventBridge Pipes针对Kinesis源内置了base64解码逻辑,只要在转换器中引用data字段(如<$.data>或<$[*].data>),就会自动将base64编码的内容解码为明文。结合JSONPath的$[*]数组投影语法,可实现对数组中所有元素的批量处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:30:43