Spring Kafka生产者偶发1分钟等待后失败问题排查求助
Spring Kafka发送事件超时1分钟问题排查求助
我在Spring Boot项目中基于Spring Kafka Template实现了事件发送功能,使用AWS MSK作为消息代理。目前遇到一个问题:部分场景下生产者发送事件会耗时1分钟才失败,虽然通过重试逻辑最终能成功发送,但这1分钟的延迟引发了很多业务问题。
生产者配置代码
@Bean public Map<String, Object> producerConfigs() throws FileNotFoundException { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers()); props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, kafkaProperties.getSecurity().getProtocol()); props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, ResourceUtils.getFile("classpath:client.truststoreks").getAbsolutePath()); props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, StringUtils.EMPTY); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); props.put(ProducerConfig.LINGER_MS_CONFIG, "100"); return props; }
生产者服务代码
public class KafkaProducerService<V> implements KafkaProducer<V> { private final KafkaTemplate<String, V> kafkaTemplate; private final KafkaTemplate<String, V> transactionLogKafkaTemplate; public KafkaProducerService(KafkaTemplate<String, V> kafkaTemplate, KafkaTemplate<String, V> transactionLogKafkaTemplate) { this.kafkaTemplate = kafkaTemplate; this.transactionLogKafkaTemplate = transactionLogKafkaTemplate; } @Override @Retryable({KafkaException.class, TimeoutException.class}) public void produce(String topic, String key, V value) { log.info("Calling producer service for producing event on topic " + topic); sendCallbackEvents(kafkaTemplate, topic, key, value); } private void sendCallbackEvents(KafkaTemplate<String, V> kafkaTemplate, String topic, String key, V value) { ProducerRecord<String, V> producerRecord = new ProducerRecord(topic, key, value); ListenableFuture<SendResult<String, V>> future = kafkaTemplate.send(producerRecord); future.addCallback(new ListenableFutureCallback<SendResult<String, V>>() { @Override public void onSuccess(SendResult<String, V> result) { log.info(String.format("Produced event to topic %s: key = %-10s value = %s", topic, key, value)); } @Override public void onFailure(Throwable ex) { log.error("Producing of data on topic {} is failed", topic, ex.getCause()); } }); } }
注:代码中已移除冗余的@Autowired注解,构造函数注入已完成依赖注入
报错信息
ERROR LogAccessor - Exception thrown when sending a message with key='xx' and payload='Event(key=value)' to topic topicName:
排查方向与解决方案建议
1. 补全生产者超时配置
当前配置未显式设置核心超时参数,默认值在AWS MSK环境下可能导致过长等待:
- 添加
REQUEST_TIMEOUT_MS_CONFIG:建议设为10000(10秒),限制单次请求超时时间 - 添加
DELIVERY_TIMEOUT_MS_CONFIG:设为30000(30秒),覆盖默认120秒的总投递超时 - 添加
MAX_BLOCK_MS_CONFIG:设为5000(5秒),限制生产者发送时的阻塞等待时间
修改后的配置片段:
props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 10000); props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 30000); props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 5000);
2. 修正SSL连接配置
当前设置SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG为空,可能引发SSL握手延迟或验证失败:
- 移除该配置,使用默认值
HTTPS,匹配AWS MSK的SSL端点验证规则 - 校验truststore文件是否包含MSK集群的CA证书,确保证书路径加载正确,避免握手时证书验证超时
3. 检查AWS MSK集群状态与网络
- 确认生产者所在VPC与MSK集群的网络连通性,检查安全组是否开放9094端口(SSL通信端口)
- 检查IAM权限(若使用IAM认证),确保生产者拥有
kafka:DescribeCluster和kafka:WriteData权限 - 查看MSK监控指标:重点关注
UnderReplicatedPartitions(分区副本不同步)、ProducerFetchResponseTime(生产者元数据请求耗时),确认集群是否负载过高
4. 优化发送逻辑与重试策略
linger.ms设为100ms会累积消息再发送,若业务对延迟敏感,可调整为0或更小值- 检查
acks配置:默认值为1,若允许一定丢包可设为0降低延迟;若需高可靠性可设为all,但会增加发送耗时 - 为重试逻辑添加指数退避,避免短时间大量重试压垮集群:
@Retryable({KafkaException.class, TimeoutException.class}, backoff = @Backoff(delay = 100, multiplier = 2, maxDelay = 1000))
5. 开启客户端DEBUG日志排查
添加日志配置到application.yml,查看SSL握手、元数据请求、消息发送的具体阶段耗时:
logging: level: org.apache.kafka: DEBUG org.springframework.kafka: DEBUG
内容的提问来源于stack exchange,提问作者Tanay Sen
相关产品推荐
相关产品推荐

