Kinesis Firehose转换时拆分单条记录为多条并存储Parquet至S3可行吗?
Kinesis Firehose拆分单条记录为多条并存储为Parquet的解决方案
原生限制说明
Kinesis Firehose的Lambda转换环节不支持直接将单条输入记录拆分为多条输出记录,核心原因是Firehose要求返回的每条记录必须携带唯一的recordId,重复的recordId会触发报错:
Multiple records were returned with the same record Id. Ensure that the Lambda function returns a unique record Id for each record.
以下是可行的替代方案,同时覆盖Redshift目标场景:
方案一:上游提前拆分记录
在数据进入Firehose之前完成拆分操作:
- 如果是自定义数据生产者,在发送数据到Firehose/Kinesis Stream前,直接将数组字段拆分为单条完整记录(比如把示例中的
key2_array元素拆成两条包含key1和对应子对象的记录),再发送。 - 如果数据来自其他数据源(如Kinesis Stream),在Stream与Firehose之间添加一个Lambda函数,消费Stream中的原始记录,拆分后生成带唯一
recordId的单条记录(比如原recordId拼接数组索引,如original-id-0、original-id-1),再转发到Firehose。
这种方式下,Firehose接收的已经是拆分完成的单条记录,可直接配置Parquet格式转换并存储到S3;若目标是Redshift,可通过Firehose直接将Parquet数据加载到Redshift。
方案二:Lambda + Kinesis Stream中间层调整架构
调整数据流转路径,用Kinesis Stream作为拆分环节的载体:
- 原始数据先发送到Kinesis Stream,而非直接到Firehose。
- 创建Lambda函数消费该Stream,对每条输入记录进行拆分:
- 遍历数组字段,为每个子对象生成一条完整记录(保留原记录的关联字段,如示例中的
key1)。 - 为每条拆分后的记录生成唯一
recordId(原recordId+UUID/数组索引均可)。
- 遍历数组字段,为每个子对象生成一条完整记录(保留原记录的关联字段,如示例中的
- 将拆分后的记录发送到Kinesis Firehose,配置Firehose进行Parquet格式转换并存储到S3;若目标是Redshift,可直接通过Firehose的Redshift加载功能,将Parquet数据导入Redshift表。
方案三:Athena批量后处理(适合非实时场景)
如果业务对实时性要求不高,可采用批量后处理方式:
- 配置Firehose将原始JSON格式数据直接存储到S3,跳过Lambda转换环节。
- 在Athena中创建外部表,映射原始JSON数据的结构。
- 使用Athena的
UNNEST函数拆分数组字段,示例SQL如下:
SELECT key1, sub.subkey1, sub.subkey2 FROM your_external_table, UNNEST(key2_array) AS sub
- 将查询结果导出为Parquet格式存储到S3;若目标是Redshift,可通过Athena的
INSERT INTO redshift_table SELECT ...语句,将拆分后的数据直接导入Redshift,或通过Redshift的COPY命令从Parquet文件加载数据。
Redshift场景补充
针对Redshift的目标需求,除上述方案外,还可直接在Redshift内部处理:
- 创建外部表指向S3上的原始JSON数据,使用
UNNEST拆分数组字段后,插入到Redshift的目标表中。 - 若采用方案一或二的拆分后数据,可通过Firehose的Redshift加载功能,直接将Parquet格式的拆分数据导入Redshift,相比JSON格式,Parquet的加载效率更高且压缩比更好。
内容的提问来源于stack exchange,提问作者tjmn
相关产品推荐
相关产品推荐

