如何使用AWS Kinesis Firehose将嵌套结构数据推送至Amazon Redshift
嗨,针对你提到的需求——把嵌套的arr字段扁平化后推送到Redshift,同时保留原有的S3完整对象存储,我来分享几个可行的实现方法,重点讲你更倾向的扁平化方案:
方案一:Firehose + Lambda 预处理(推荐,符合你的扁平化需求)
这是最直接的端到端处理方式,利用Firehose的Lambda转换功能在数据到达Redshift前完成扁平化:
编写Lambda转换函数
你需要写一个Python/Node.js的Lambda函数,接收Firehose批量传递的原始数据,遍历每个JSON对象,将外层的field1、field2与arr数组中的每个元素组合成新的扁平化对象,最后输出这些处理后的记录。举个Python示例(处理批量数据):
import json import base64 def lambda_handler(event, context): output_records = [] for record in event['records']: # 解码Firehose传递的base64格式数据 payload = json.loads(base64.b64decode(record['data'])) # 遍历数组生成扁平化记录 for inner_item in payload['arr']: flattened_record = { 'field1': payload['field1'], 'field2': payload['field2'], 'inner_field1': inner_item['inner_field1'], 'inner_field2': inner_item['inner_field2'] } # 编码回Firehose要求的格式 output_record = { 'recordId': record['recordId'], 'result': 'Ok', 'data': base64.b64encode(json.dumps(flattened_record).encode('utf-8')).decode('utf-8') } output_records.append(output_record) return {'records': output_records}配置Firehose启用Lambda转换
在Firehose控制台的配置页面,找到「数据转换」环节,启用Lambda转换并选择你创建的函数。同时调整Redshift目标配置,确保目标表的列(field1、field2、inner_field1、inner_field2)和扁平化后的字段一一对应,这样处理后的数据就能直接写入目标表。
方案二:Redshift端 COPY + UNNEST 后处理
如果不想修改Firehose的数据流,也可以先把完整的JSON数据(包含arr字段)推送到Redshift,再在Redshift内部完成扁平化:
调整Firehose配置
修改Firehose的Redshift目标设置,确保推送完整的JSON对象(而非仅field1、field2)到Redshift的临时表(可用SUPER类型存储整个JSON,或用VARCHAR类型)。用UNNEST展开数组并插入目标表
先创建临时表存储原始完整数据,再通过UNNEST函数展开数组,关联外层字段插入到最终目标表:-- 创建临时表存储原始完整JSON CREATE TEMP TABLE raw_temp_data (raw_json SUPER); -- 从S3 COPY数据(Firehose会把完整JSON同步到S3,再推送到Redshift) COPY raw_temp_data FROM 's3://your-firehose-bucket/path/' IAM_ROLE 'arn:aws:iam::your-account-id:role/your-redshift-role' FORMAT JSON 'auto'; -- 扁平化数据插入目标表 INSERT INTO your_target_table (field1, field2, inner_field1, inner_field2) SELECT raw_json.field1, raw_json.field2, inner_obj.inner_field1, inner_obj.inner_field2 FROM raw_temp_data, UNNEST(raw_json.arr) AS t(inner_obj);你可以把这段SQL做成定时任务(比如用Redshift定时查询或AWS Glue),自动完成数据的扁平化转换。
补充:关于SUPER类型的实现
虽然你一开始没找到相关文档,但Firehose确实支持推送SUPER类型数据到Redshift:
- 在Redshift目标表中创建
arr_col SUPER类型的列; - 配置Firehose的Redshift COPY选项为
FORMAT JSON 'auto',Redshift会自动将嵌套的JSON数组解析为SUPER类型的数组; - 之后你可以直接在Redshift中用
UNNEST(arr_col)查询或转换数据,也可以保留嵌套结构用于灵活分析。不过如果核心需求是生成扁平化表,前面两种方案会更直接。
注意事项
- Lambda转换时要关注性能:大流量场景下,调整Lambda的内存配置(建议至少1024MB)和超时时间,避免处理延迟;
- Redshift端处理时,注意临时表的生命周期管理,避免占用过多存储空间;
- 测试阶段先用小批量数据验证处理逻辑,确保字段映射正确。
内容的提问来源于stack exchange,提问作者Damien

