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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 03:45:34