使用ReplyingKafkaTemplate实现Kafka同步通信时Reply Topic消费滞后问题
问题背景
产品包含多个微服务,部分业务场景下微服务TryServiceOne需要将请求转发至另一微服务TryServiceThree处理,终端用户需等待API返回响应,因此采用ReplyingKafkaTemplate实现双向同步通信,即时向调用方返回结果。目前业务功能运行正常,但监控发现REPLY Topic存在消费滞后(LAG),导致告警系统频繁触发告警;但实际RequestReplyFuture已正常读取并成功处理消息,Kafka broker侧统计的消费滞后量仍持续上涨,需要给出规避该消费滞后问题的方案。
重要说明:采用多节点集群模式部署微服务,因此通过自定义分区策略,将响应/回复Topic的所有消息固定路由至同一个分区。
现有配置代码
TryServiceOne 配置
KafkaConfiguration.class 代码
@Bean public Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(org.apache.kafka.clients.producer.ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers); props.put(org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); return props; } @Bean public Map<String,Object> consumerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId); return props; } @Bean public ProducerFactory<String, RequestModel> requestProducerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } @Bean public KafkaTemplate<String, RequestModel> kafkaTemplate() { return new KafkaTemplate<>(requestProducerFactory()); } @Bean public ReplyingKafkaTemplate<String, RequestModel, ResponseModel> replyKafkaTemplate(ProducerFactory<String, RequestModel> pf, KafkaMessageListenerContainer<String, ResponseModel> container){ return new ReplyingKafkaTemplate<>(pf, container); } @Bean public KafkaMessageListenerContainer<String, ResponseModel> replyContainer(ConsumerFactory<String, ResponseModel> cf) { TopicPartitionOffset topicPartitionOffset = new TopicPartitionOffset("RESPONSE_TOPIC",0); ContainerProperties containerProperties = new ContainerProperties(topicPartitionOffset); containerProperties.setAckMode(ContainerProperties.AckMode.MANUAL); return new KafkaMessageListenerContainer<>(cf, containerProperties); }
SendAndReceive服务组件实现代码
RequestModel requestModel= new RequestModel(); distributorRequestEvent.setDistributorModel(producerRecord); // 构造生产者消息 ProducerRecord<String, RequestModel> record = new ProducerRecord<String, RequestModel>("REQUEST_TOPIC", requestModel); // 在消息头设置回复Topic record.headers().add(new RecordHeader(KafkaHeaders.REPLY_TOPIC, "RESPONSE_TOPIC".getBytes(StandardCharsets.UTF_8))); kafkaTemplate.setDefaultReplyTimeout(Duration.ofSeconds(30)); LOGGER.info("Sending message ... {}",producerRecord); RequestReplyFuture<String, RequestModel, ResponseModel> sendAndReceive = kafkaTemplate.sendAndReceive(record); // 确认消息生产成功 SendResult<String, RequestModel> sendResult = sendAndReceive.getSendFuture().get(); // 获取消费到的响应记录 ConsumerRecord<String, ResponseModel> consumerRecord = sendAndReceive.get(); return consumerRecord.value();
TryServiceThree 微服务配置
Kafka 基础配置
@Bean public Map<String, Object> consumerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId); props.put(JsonDeserializer.TYPE_MAPPINGS,RequestModel.class); return props; } @Bean public Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG,CustomPartitioner.class); props.put(ProducerConfig.ACKS_CONFIG, "all"); return props; } @Bean public ConsumerFactory<String, RequestModel> requestConsumerFactory() { JsonDeserializer<RequestModel> deserializer = new JsonDeserializer<>(RequestModel.class); deserializer.setRemoveTypeHeaders(false); deserializer.addTrustedPackages("*"); deserializer.setUseTypeMapperForKey(true); return new DefaultKafkaConsumerFactory<>(consumerConfigs(), new StringDeserializer(), deserializer); } @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, RequestModel>> requestListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, RequestModel> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(requestConsumerFactory()); factory.setReplyTemplate(replyTemplate()); return factory; } @Bean public ProducerFactory<String, ResponseModel> replyProducerFactory() { ProducerFactory<String, ResponseModel> producerFactory = new DefaultKafkaProducerFactory<>(producerConfigs()); return producerFactory; } @Bean public KafkaTemplate<String, ResponseModel> replyTemplate() { return new KafkaTemplate<>(replyProducerFactory()); }
自定义分区器实现代码
public class CustomPartitioner implements Partitioner { @Override public int partition(String s, Object o, byte[] bytes, Object o1, byte[] bytes1, Cluster cluster) { return 0; } @Override public void close() { } @Override public void configure(Map<String, ?> map) { } }
问题根因与解决方案
根因分析
消费滞后持续上涨的核心原因是当前配置中回复监听器被设置为MANUAL手动提交偏移量模式,但代码中无任何手动提交offset的逻辑。ReplyingKafkaTemplate内置的回复消费者仅在内存中匹配对应correlation id的响应消息交给业务线程处理,不会自动提交消费位移,Broker侧长期未收到该消费组的offset提交请求,会判定所有已拉取的消息都未被消费,LAG数值自然持续累加。
另外多节点部署场景下所有响应被固定路由到分区0,所有TryServiceOne实例会竞争同一个分区的消费权,同一时间仅一个实例能消费到该分区消息,其余实例空跑,也会放大offset提交异常的影响。
修复方案
- 方案一(最简便,无侵入):直接把回复容器的AckMode改成
AckMode.RECORD,容器每处理完一条响应消息就会自动提交offset,不需要额外写提交逻辑,完全适配ReplyingKafkaTemplate的场景,不会影响业务正常匹配响应。
修改代码如下:// 替换原有的MANUAL手动提交模式 containerProperties.setAckMode(ContainerProperties.AckMode.RECORD); - 方案二(保留手动提交模式):如果业务必须用手动提交,需要给
replyContainer配置AcknowledgingMessageListener,在消息被ReplyingKafkaTemplate正常接收后手动调用acknowledge()方法提交位移,注意不要覆盖ReplyingKafkaTemplate自带的消息匹配逻辑,需先将消息转给模板内部处理再提交。 - 方案三(优化集群部署问题):不要把所有响应都固定发到分区0,建议给RESPONSE_TOPIC扩容分区数,分区数和TryServiceOne的节点数对齐,自定义分区器按请求头里的
KafkaHeaders.CORRELATION_ID做哈希路由,让每个TryServiceOne节点消费固定的分区,既避免消费重平衡开销,也能从架构上避免单分区消费瓶颈。 - 额外配置优化:给回复消费者配置
ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG为false,避免自动提交和容器提交模式冲突;同时给消费组配置合理的auto.offset.reset策略,防止异常重启后从最早位置消费导致大量无效消息拉取。
内容的提问来源于stack exchange,提问作者Bharat Chary
相关产品推荐
相关产品推荐

