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

如何使用AWS Kinesis Firehose将嵌套结构数据推送至Amazon Redshift

实现Kinesis Firehose嵌套数据扁平化入Redshift的方案

嗨,针对你提到的需求——把嵌套的arr字段扁平化后推送到Redshift,同时保留原有的S3完整对象存储,我来分享几个可行的实现方法,重点讲你更倾向的扁平化方案:

方案一:Firehose + Lambda 预处理(推荐,符合你的扁平化需求)

这是最直接的端到端处理方式,利用Firehose的Lambda转换功能在数据到达Redshift前完成扁平化:

  1. 编写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}
    
  2. 配置Firehose启用Lambda转换
    在Firehose控制台的配置页面,找到「数据转换」环节,启用Lambda转换并选择你创建的函数。同时调整Redshift目标配置,确保目标表的列(field1、field2、inner_field1、inner_field2)和扁平化后的字段一一对应,这样处理后的数据就能直接写入目标表。

方案二:Redshift端 COPY + UNNEST 后处理

如果不想修改Firehose的数据流,也可以先把完整的JSON数据(包含arr字段)推送到Redshift,再在Redshift内部完成扁平化:

  1. 调整Firehose配置
    修改Firehose的Redshift目标设置,确保推送完整的JSON对象(而非仅field1、field2)到Redshift的临时表(可用SUPER类型存储整个JSON,或用VARCHAR类型)。

  2. 用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:

  1. 在Redshift目标表中创建arr_col SUPER类型的列;
  2. 配置Firehose的Redshift COPY选项为FORMAT JSON 'auto',Redshift会自动将嵌套的JSON数组解析为SUPER类型的数组;
  3. 之后你可以直接在Redshift中用UNNEST(arr_col)查询或转换数据,也可以保留嵌套结构用于灵活分析。不过如果核心需求是生成扁平化表,前面两种方案会更直接。

注意事项

  • Lambda转换时要关注性能:大流量场景下,调整Lambda的内存配置(建议至少1024MB)和超时时间,避免处理延迟;
  • Redshift端处理时,注意临时表的生命周期管理,避免占用过多存储空间;
  • 测试阶段先用小批量数据验证处理逻辑,确保字段映射正确。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 15:42:37