【发布时间】: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
我们正在使用 Spring Cloud Streams Hoxton.SR4 来使用来自 Kafka 主题的消息。我们启用了 spring.cloud.stream.bindings..consumer.batch-mode=true,每次轮询获取 2000 条记录。我想知道是否有一种方法可以手动确认/提交整个批次。
【问题讨论】:
标签: spring-kafka spring-cloud-stream
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=''}}
【讨论】:
有没有办法在批处理消费者中检索消息列表,包括它们的标头,类似的东西(更新了上面的示例)
@Bean
Consumer<List<Message<Foo>>> consume() {
return list -> {
list.forEach(msg -> {
System.out.println(msg.getPayload());
});
};
}
虽然我的问题是关于带有标头的批量消费(如 KafkaHeaders.MESSAGE_KEY),因为我的生产者部分在 Key 中发送所需的数据,而在 Payload 中休息。我没有找到与我要查找的内容最接近的主题。
【讨论】: