没有待答复:ConsumerRecord

Ken*_*hia 4 java apache-kafka spring-kafka

我正在尝试使用 ReplyingKafkaTemplate,并且间歇性地不断看到下面的消息。

没有待处理的回复:ConsumerRecord(主题 = 请求回复主题,分区 = 8,偏移量 = 1,CreateTime = 1544653843269,序列化键大小 = -1,序列化值大小 = 1609,标头 = RecordHeaders(标头 = [RecordHeader(键 = kafka_correlationId,值 = [-14, 65, 21, -118, 70, -94, 72, 87, -113, -91, 92, 72, -124, -110, -64, -94])], isReadOnly = false),key = null,correlationId: [-18271255759235816475365319231847350110],可能超时,或使用共享回复主题

它将源于下面的代码

RequestReplyFuture<K, V, R> future = this.futures.remove(correlationId);
if (future == null) {
  if (this.sharedReplyTopic) {
    if (this.logger.isDebugEnabled()) {
      this.logger.debug(missingCorrelationLogMessage(record, correlationId));
    }
  }
  else if (this.logger.isErrorEnabled()) {
    this.logger.error(missingCorrelationLogMessage(record, correlationId));
  }
}
Run Code Online (Sandbox Code Playgroud)

但只是间歇性地发生

我还将共享的replyTopic设置为 false ,如下所示,并尝试强制更长的超时

ReplyingKafkaTemplate<String, Object, Object> replyKafkaTemplate = new ReplyingKafkaTemplate<>(pf, container);
        replyKafkaTemplate.setSharedReplyTopic(false);
        replyKafkaTemplate.setReplyTimeout(10000);
        return replyKafkaTemplate;
Run Code Online (Sandbox Code Playgroud)

我的容器如下

@Bean
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, Object>> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());

    factory.setBatchListener(false);
    factory.getContainerProperties().setPollTimeout(1000);
    factory.getContainerProperties().setIdleEventInterval(10000L);
    factory.setConcurrency(3);
    factory.setReplyTemplate(kafkaTemplate());
    return factory;
}
Run Code Online (Sandbox Code Playgroud)

Gar*_*ell 5

如果是断断续续的,则很可能回复时间过长才到达。消息看起来很清楚

可能超时,或者使用共享回复主题

每个客户端实例必须使用自己的回复主题或专用分区。

编辑

如果收到的消息的相关 ID 与 this.futures 中当前的条目(待回复)不匹配,您将收到日志。仅在以下情况下才会发生这种情况:

  1. 请求超时(这种情况会有相应的WARN日志)。
  2. 模板被 stop() 处理(在这种情况下 this.futures 被清除)。
  3. 已处理的回复由于某种原因被重新传递(不应该发生)。
  4. 在将密钥添加到 this.futures 之前收到回复(不会发生,因为它是在 send() 记录之前插入的)。
  5. 服务器端对同一请求发送 2 个或多个回复。
  6. 其他一些应用程序正在将数据发送到同一回复主题。如果您可以使用 DEBUG 日志记录来重现它,这将会有所帮助,因为这样我们也会在发送时记录相关键。