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

已启用GZIP压缩,Kinesis Firehose仍向S3投递未压缩文件

问题描述

我通过Lambda函数将JSON字符串直接写入Kinesis Firehose流,期望将批量记录以压缩的.gz文件形式投递至S3。
尽管已将流的“Destination settings > Compression for data records”设置为GZIP,但投递的文件虽带有.gz后缀,实际却是明文格式。可通过两点验证:

  • 下载文件后无需修改即可直接以文本形式打开;
  • 执行gzip -d ~/path/my_file.gz命令时返回gzip: /path/my_file.gz: not in gzip format。

为何启用压缩后Firehose仍投递未压缩的数据?我是否遗漏了某些配置?

代码示例

Lambda 代码

import json
import boto3
firehose = boto3.client("firehose")

record = {'field_1': 'test'}               # dict/json
record_string = json.dumps(record) + '\n'  # Firehose expects ndjson

response = firehose.put_record(
    DeliveryStreamName=my_stream_name,
    Record={ 'Data': record_string }
)

Firehose Terraform 配置

resource "aws_kinesis_firehose_delivery_stream" "my_firehose_stream" {
  name        = my_stream_name
  destination = "extended_s3"

  extended_s3_configuration {
    role_arn   = my_role_arn
    bucket_arn = my_bucket_arn

    prefix              = "my_prefix/!{partitionKeyFromQuery:extracted}/"
    error_output_prefix = "my_error_prefix/"

    buffering_size      = 64     # MB
    buffering_interval  = 900    # seconds
    compression_format  = "GZIP" # Compress as GZIP

    # Enabled to dynamic extract
    processing_configuration {
      enabled = true
      processors {
        type = "MetadataExtraction"
        parameters {
          parameter_name  = "JsonParsingEngine"
          parameter_value = "JQ-1.6"
        }
        parameters {
          parameter_name  = "MetadataExtractionQuery"
          parameter_value = "{extracted:.extracted}"
        }
      }
    }

    dynamic_partitioning_configuration {
      enabled        = true
    }
  }
}

问题原因及解决方案

问题出在MetadataExtraction处理器的配置上:当启用MetadataExtraction并开启动态分区时,如果Firehose无法从记录中提取到指定的元数据(比如你的示例中写入的record里根本没有extracted字段),Firehose会将这条记录路由到错误输出路径;如果所有记录都触发错误,Firehose会跳过压缩步骤,直接将明文数据以.gz后缀写入S3,这是Firehose的特殊处理逻辑。

你的Lambda写入的{'field_1': 'test'}没有extracted字段,导致MetadataExtraction处理器执行失败,进而触发了跳过压缩的行为。

解决步骤:

  • 修正元数据提取逻辑:要么修改Lambda写入包含extracted字段的数据,要么调整MetadataExtractionQuery为现有字段,比如{extracted:.field_1}。
  • 验证错误记录:检查Firehose的错误输出路径(my_error_prefix/),确认是否有大量错误记录,以此验证是否是处理器执行失败导致的压缩跳过。
  • 临时验证:可以先关闭processing_configuration,测试压缩功能是否正常,确认问题根源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 21:33:21