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

Kafka MirrorMaker2(Strimzi)生产者配置位置咨询

Kafka MirrorMaker2(Strimzi)生产者配置正确位置说明

错误原因分析

你看到的flushing 197026 outstanding messages for offset commit日志,说明MirrorSourceConnector拉取了大量消息但未成功发送到目标集群,导致偏移量提交前堆积了大量未处理消息。核心问题是生产者发送效率不足,而你之前的配置位置或参数格式错误,导致生产者配置未生效。

正确的生产者配置位置

在Strimzi的KafkaMirrorMaker2 CR中,生产者配置需要根据作用范围放在以下两个位置:

1. 针对单个Mirror的SourceConnector(最常用)

如果只想给复制业务的生产者(即从源集群拉取后发往目标集群的生产者)配置参数,需要放在mirrors[].sourceConnector.config下,直接使用Kafka生产者原生参数名,不要加producer.前缀。

示例配置:

mirrors:
  - sourceCluster: source-kafka
    targetCluster: target-kafka
    # ... 其他Mirror配置
    sourceConnector:
      config:
        # ... 其他SourceConnector配置
        batch.size: 65536  # 直接用原生参数名,无需加producer.前缀
        buffer.memory: 52428800  # 50MB,单位为字节
        linger.ms: 50  # 可选配置,提升批量发送效率

2. 实例级全局生产者配置

如果要给当前MM2实例下的所有Connector(包括SourceConnector、HeartbeatConnector、CheckpointConnector)统一配置生产者参数,需要放在instances[].producer.config下,同样使用原生参数名,不要加producer.前缀:

instances:
  - enabled: true
    name: mirror-maker
    # ... 其他实例配置
    producer:
      config:
        batch.size: 65536
        buffer.memory: 52428800

你当前配置的问题

  • spec.config下的参数属于Connect Worker全局配置,不属于生产者参数,无法作用到消息发送的生产者。
  • instances[].producer.config里的producer.batch.size、producer.buffer.memory带了多余的producer.前缀,Strimzi无法识别这些参数。
  • mirrors[]下直接配置的offset.flush.timeout.ms等是Mirror通用配置,不是生产者参数,无法影响发送逻辑。
  • sourceConnector.config里的producer.batch.size同样带了前缀,应改为batch.size。
  • 额外注意:你配置的max.poll.record参数名错误,正确名称为max.poll.records,该参数控制消费者单次拉取消息数,也会影响消息堆积情况,建议修正。

修正后的关键配置片段示例

kafkaMirrorMaker2:
 spec:
  instances:
   - enabled: true
     name: mirror-maker
     replicas: 3
     enableMetrics: true
     targetCluster: target-kafka
     # 全局生产者配置(可选)
     producer:
       config:
         batch.size: 65536
         buffer.memory: 52428800
     clusters:
       source-kafka:
         bootstrapServers: sourcekafka:9092
       target-kafka:
         bootstrapServers: targetkafka:9092
     mirrors:
       - sourceCluster: source-kafka
         targetCluster: target-kafka
         topicsPattern: "mytopic"
         groupsPattern: "my-replicator"
         sourceConnector:
           config:
             max.poll.records: 2000  # 修正参数名
             tasks.max: 9
             consumer.auto.offset.reset: latest
             offset.flush.timeout.ms: 3000
             consumer.request.timeout.ms: 15000
             # 针对SourceConnector的生产者配置
             batch.size: 65536
             buffer.memory: 52428800
             linger.ms: 50
         # ... 其他Connector配置
     resources:
       requests:
         memory: 2Gi
         cpu: 3
       limits:
         memory: 4Gi
         cpu: 3

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 08:55:17