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、目标分区等信息,需要调整回调的传参逻辑,具体实现方式如下:
- 自定义回调类时,不要仅传递payload,而是将完整的
ProducerRecord实例作为成员变量存入回调对象,异常触发时可以直接从实例中提取需要的key、value、分区、topic信息。 - 不要在回调线程中直接同步调用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
相关产品推荐
相关产品推荐

