0
votes

We are using spring cloud streams Hoxton.SR4 to consume messages from Kafka topic. We've enabled spring.cloud.stream.bindings..consumer.batch-mode=true, fetching 2000 records per poll. I would like to know if there is a way we can manually acknowledge/commit entire batch.

1

1 Answers

0
votes

SR4 is quite old; the current Hoxton release is SR9 and the current spring cloud stream version is 3.0.10.RELEASE (Hoxton.SR9 pulls in 3.0.9).

You need to consume a Message and get the acknowledgment from a header.

@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 + "]";
        }

    }

}

Properties for Boot 2.3.6 and 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

Properties for Boot 2.4.0 and 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=''}}