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

使用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解析引擎版本。

注意:请确保写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 06:39:32