【问题标题】:Spring Cloud Stream Kafka batch - Manual Commit Entire batchSpring Cloud Stream Kafka 批处理 - 手动提交整个批处理
【发布时间】:2020-12-10 05:21:15
【问题描述】:

我们正在使用 Spring Cloud Streams Hoxton.SR4 来使用来自 Kafka 主题的消息。我们启用了 spring.cloud.stream.bindings..consumer.batch-mode=true,每次轮询获取 2000 条记录。我想知道是否有一种方法可以手动确认/提交整个批次。

【问题讨论】:

    标签: spring-kafka spring-cloud-stream


    【解决方案1】:

    SR4 已经很老了;当前 Hoxton 版本为 SR9,当前 Spring Cloud Stream 版本为 3.0.10.RELEASE(Hoxton.SR9 引入 3.0.9)。

    您需要使用 Message 并从标头中获取确认。

    @SpringBootApplication
    public class So652289261Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So652289261Application.class, args);
        }
    
        @Bean
        Consumer<Message<List<Foo>>> consume() {
            return msg -> {
                System.out.println(msg.getPayload());
                msg.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class).acknowledge();
            };
        }
    
        @Bean
        public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer() {
            return (container, dest, group) -> container.getContainerProperties()
                    .setCommitLogLevel(LogIfLevelEnabled.Level.INFO);
        }
    
        @Bean
        public ApplicationRunner runner(KafkaTemplate<byte[], byte[]> template) {
            return args -> {
                template.send("consume-in-0", "{\"bar\":\"baz\"}".getBytes());
                template.send("consume-in-0", "{\"bar\":\"qux\"}".getBytes());
            };
        }
    
        public static class Foo {
    
            private String bar;
    
            public Foo() {
            }
    
            public Foo(String bar) {
                this.bar = bar;
            }
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    Boot 2.3.6 和 Cloud Hoxton.SR9 的属性

    spring.cloud.stream.bindings.consume-in-0.group=so65228926
    spring.cloud.stream.bindings.consume-in-0.consumer.batch-mode=true
    spring.cloud.stream.kafka.bindings.consume-in-0.consumer.auto-commit-offset=false
    
    spring.kafka.producer.properties.linger.ms=50
    

    Boot 2.4.0 和 Cloud 2020.0.0-M6 的属性

    spring.cloud.stream.bindings.consume-in-0.group=so65228926
    spring.cloud.stream.bindings.consume-in-0.consumer.batch-mode=true
    spring.cloud.stream.kafka.bindings.consume-in-0.consumer.ack-mode=MANUAL
    
    spring.kafka.producer.properties.linger.ms=50
    
    [Foo [bar=baz], Foo [bar=qux]]
    ... Committing: {consume-in-0-0=OffsetAndMetadata{offset=14, leaderEpoch=null, metadata=''}}
    

    【讨论】:

      【解决方案2】:

      有没有办法在批处理消费者中检索消息列表,包括它们的标头,类似的东西(更新了上面的示例)

      @Bean
      Consumer<List<Message<Foo>>> consume() {
          return list -> {
              list.forEach(msg -> {
                  System.out.println(msg.getPayload());
              });
          };
      }
      

      虽然我的问题是关于带有标头的批量消费(如 KafkaHeaders.MESSAGE_KEY),因为我的生产者部分在 Key 中发送所需的数据,而在 Payload 中休息。我没有找到与我要查找的内容最接近的主题。

      【讨论】:

      • 另外,我尝试通过设置 spring.cloud.stream.bindings..consumer.headerMode: embeddedHeaders 并期望标头是 FooWithHeaders 对象的一部分,但它们为空跨度>
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-10
      • 1970-01-01
      相关资源
      最近更新 更多