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
相关产品推荐
相关产品推荐

