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

