使用Python CDK创建启用动态分区的Kinesis Firehose传输流
报错根因
Kinesis Firehose动态分区能力依赖托管的元数据提取处理器从记录中解析分区键,不允许单独开启动态分区而关闭处理配置,这是你收到Processing Configuration is not enabled when DataPartitioning is enabled报错的直接原因。
核心配置规则
processing_configuration块是开启动态分区的必填项,各字段按照固定规则填写即可,无需自定义业务逻辑:
enabled:必须设为True,你查到的官方参考示例中enabled=False是关闭数据处理的通用配置,不适用于动态分区场景。- 处理器
type:固定填MetadataExtraction,这是Firehose提供的全托管内置处理器,专门用于从记录Payload中提取分区键,无需额外部署Lambda函数。 - 处理器
parameters需传入两个必填参数:- 参数1:
parameter_name设为MetadataExtractionQuery,parameter_value设为{log_type:.log_type},采用JQ语法定义提取规则:从每条JSON日志的根路径下读取log_type字段值,将其映射为名为log_type的分区键,和你S3前缀里的!{partitionKeyFromQuery:log_type}占位符对应。 - 参数2:
parameter_name设为JsonParsingEngine,parameter_value设为JQ-1.6,指定使用的JSON解析引擎版本。
- 参数1:
注意:请确保写入Firehose的日志为标准JSON格式,且每条记录根节点携带
log_type字段,无法正常提取分区键的异常记录会被自动写入你配置的error_output_prefix对应路径。
修正后完整代码
analytics_delivery_stream = kinesisfirehose.CfnDeliveryStream( self, "AnalyticsDeliveryStream", delivery_stream_name='analytics', extended_s3_destination_configuration=kinesisfirehose.CfnDeliveryStream.ExtendedS3DestinationConfigurationProperty( bucket_arn=f'arn:aws:s3:::{analytic_bucket_name}', buffering_hints=kinesisfirehose.CfnDeliveryStream.BufferingHintsProperty( interval_in_seconds=60 ), dynamic_partitioning_configuration = kinesisfirehose.CfnDeliveryStream.DynamicPartitioningConfigurationProperty( enabled=True, retry_options=kinesisfirehose.CfnDeliveryStream.RetryOptionsProperty( duration_in_seconds=123 ) ), # 新增必填的处理配置,用于提取log_type分区键 processing_configuration=kinesisfirehose.CfnDeliveryStream.ProcessingConfigurationProperty( enabled=True, processors=[ kinesisfirehose.CfnDeliveryStream.ProcessorProperty( type="MetadataExtraction", parameters=[ kinesisfirehose.CfnDeliveryStream.ProcessorParameterProperty( parameter_name="MetadataExtractionQuery", parameter_value="{log_type:.log_type}" ), kinesisfirehose.CfnDeliveryStream.ProcessorParameterProperty( parameter_name="JsonParsingEngine", parameter_value="JQ-1.6" ) ] ) ] ), compression_format="UNCOMPRESSED", role_arn=firehose_role.role_arn, prefix="!{partitionKeyFromQuery:log_type}/!{timestamp:yyyy}/!{timestamp:MM}/!{timestamp:dd}/", error_output_prefix="errors/!{firehose:error-output-type}/!{timestamp:yyyy}/anyMonth/!{timestamp:dd}/", ) )
效果说明
部署上述配置后,Firehose会自动按log_type字段值分区:
log_type为type_A_log的记录会存入S3的type_A_log/年/月/日/目录log_type为type_B_log的记录会存入S3的type_B_log/年/月/日/目录
完全匹配你的分区需求,且MetadataExtraction为托管处理器,不需要给Firehose绑定的角色额外添加自定义资源调用权限。
内容的提问来源于stack exchange,提问作者Naxi
相关产品推荐
相关产品推荐

