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

Kafka Producer报记录过期超时的配置优化与消息重发咨询

问题解答

针对配置遗漏导致过期报错的调整方案

Expiring X record(s) for xxx: 120004 ms has passed since batch creation报错的触发逻辑是:消息从批次创建开始计算,总耗时超过生产者默认的delivery.timeout.ms阈值(默认值120000ms,和报错中的120004ms完全匹配),消息被生产者强制丢弃。结合当前配置和后续负载升高的预期,需要补充调整以下生产者配置:

  • delivery.timeout.ms:建议调整为300000ms(5分钟),该值为消息从进入生产者缓冲区到最终发送成功/失败的总超时上限,注意配置值必须大于linger.ms与request.timeout.ms之和,否则启动时会报配置错误。
  • request.timeout.ms:默认值30000ms(30秒),是生产者单次向broker发送请求后等待响应的超时阈值,负载升高后broker副本同步、磁盘刷盘延迟会上升,建议调整为60000ms(1分钟),减少不必要的单次请求超时。
  • enable.idempotence=true:当前配置了acks=all,开启幂等生产者后会自动调整重试、在途请求数相关参数,避免重试导致的消息重复、乱序问题;同时配套调整retry.backoff.ms=500,默认100ms的重试间隔在高负载下会频繁发起无效重试,加大broker压力。
  • buffer.memory:默认值33554432(32MB),是生产者用来缓存待发送消息的总内存大小,如果后续生产速率超过broker接收速率,缓冲区打满会导致消息阻塞无法发送,最终触发过期,建议根据吞吐量调整为67108864(64MB)~134217728(128MB)。
  • 额外注意:配置调整后仍需监控对应分区所在broker的磁盘IO、ISR列表状态,acks=all模式下需要所有ISR副本确认写入,如果ISR频繁收缩、broker磁盘IO打满,仅调整生产者配置无法彻底解决发送超时问题。

过期消息的提取与重发实现方案

当前回调实现仅存储了消息payload,无法在报错时获取完整的消息key、目标分区等信息,需要调整回调的传参逻辑,具体实现方式如下:

  1. 自定义回调类时,不要仅传递payload,而是将完整的ProducerRecord实例作为成员变量存入回调对象,异常触发时可以直接从实例中提取需要的key、value、分区、topic信息。
  2. 不要在回调线程中直接同步调用send方法重发,回调逻辑运行在生产者的IO线程上,同步重发会阻塞IO线程,加剧发送阻塞,正确做法是将过期消息投递到业务侧独立的重试队列,由专门的工作线程异步执行退避重试。

参考实现代码:

// 发送时传入完整ProducerRecord实例,不要只传payload
ProducerRecord<String, String> rec = new ProducerRecord<>("myTopic", 1, "myKey", "json-payload-here");
producer.send(rec, new ProducerCallback(rec));

private class ProducerCallback implements Callback {
  private static final String _ME = "onCompletion";
  private final ProducerRecord<String, String> failedRecord;

  public ProducerCallback(ProducerRecord<String, String> record) {
    this.failedRecord = record;
  }

  @Override
  public void onCompletion(RecordMetadata recordMetadata, Exception e) {
    if (e == null) {
      LOG.logp(Level.FINEST, _CL, _ME, "Published kafka event, offset: " + recordMetadata.offset());
      return;
    }
    LOG.log(Level.SEVERE, "Publish failed, " + e.getMessage(), e);
    // 判定为过期超时的消息
    if (e instanceof TimeoutException) {
      String topic = failedRecord.topic();
      Integer partition = failedRecord.partition();
      String key = failedRecord.key();
      String payload = failedRecord.value();
      // 投递到业务侧独立重试队列,由工作线程异步退避重试,不要在此处直接调用producer.send
      retryQueue.offer(new RetryEvent(topic, partition, key, payload, LocalDateTime.now()));
    }
  }
}

如果对消息可靠性要求高,可以将待重试的消息持久化到本地磁盘(如本地日志、嵌入式KV存储),避免生产者进程崩溃导致待重试消息丢失。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 10:09:17