【问题标题】:Spring Cloud @StreamListener doesn't provide an Acknowledgement header even auto-commit set to falseSpring Cloud @StreamListener 不提供确认标头,即使自动提交设置为 false
【发布时间】:2018-11-01 01:08:51
【问题描述】:

我目前陷入了手动管理 Kafka 偏移和提交的简单示例。我有一个带有 Spring Cloud Streams 的应用程序,它设置了 enable.auto.commit = false(在打印 ConsumerValues 时在启动日志中看到),但是当我解析消息时它没有提供确认头。

这是我的听众:

@StreamListener(Sink.INPUT)
public void handleSchedulerMessage(@Payload SchedulerEvent event, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) {
    log.debug("[message={}]", event);
    // todo: processing
    log.debug("Event processed successfully [event={}]", event);
}

用于配置的 YAML 也很简单:

spring:
  application:
    name: scheduler
  cloud:
    stream:
      kafka:
        binder:
          brokers: *kafka-broker*:9092
          zkNodes: *zookeeper*:2181
      bindings:
        input:
          destination: scheduler
          contentType: application/json
          consumer:
            autoCommitOffset: false

而且当我发送消息时,立即弹出错误:

2018-05-22 11:38:32.470 ERROR 11651 --- [container-0-C-1] o.s.integration.handler.LoggingHandler   : org.springframework.messaging.MessageHandlingException: Missing header 'kafka_acknowledgment' for method parameter type [interface org.springframework.kafka.support.Acknowledgment], failedMessage=GenericMessage [payload=byte[38], headers={kafka_offset=9, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@fdac355, deliveryAttempt=3, kafka_timestampType=CREATE_TIME, kafka_receivedMessageKey=null, kafka_receivedPartitionId=0, kafka_receivedTopic=scheduler, kafka_receivedTimestamp=1526981909241, contentType=application/json}]

收到的消息不包含所需的标头,这与禁用 autoCommit 时文档所说的相反:

Whether to autocommit offsets when a message has been processed. If set to false, a header with the key kafka_acknowledgment of the type org.springframework.kafka.support.Acknowledgment header will be present in the inbound message. Applications may use this header for acknowledging messages.

代码并不复杂,我没有使用任何预先生成的项目。示例并不能解释我所做的更多,所以我不知道我可能会错过什么。

【问题讨论】:

    标签: spring-boot apache-kafka spring-cloud spring-cloud-stream spring-kafka


    【解决方案1】:

    看起来您在 YAML 中丢失了一级缩进。

    根据文档,属性必须是这样的:

    spring.cloud.stream.kafka.bindings.input.consumer.autoCommitOffset=false
    

    但是你的样本是这样的:

    spring.cloud.stream.bindings.input.consumer.autoCommitOffset=false
    

    注意中间多余的.kafka.

    我不知道如何帮助正确管理 YAML,但我们必须这样做才能使其正常工作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-10-01
      • 1970-01-01
      • 1970-01-01
      • 2017-03-11
      • 1970-01-01
      • 1970-01-01
      • 2020-08-31
      • 2016-04-04
      相关资源
      最近更新 更多