【发布时间】:2022-10-17 20:56:57
【问题描述】:
我在 Spring Boot 应用程序中使用 Spring Kafka。我正在尝试使用 Kafka ConsumerInterceptor 在提交偏移量时进行拦截。
这似乎工作生产者事务未启用,但事务已打开 Interceptor::onCommit 不再被调用。
以下最小示例一切都按预期工作:
@SpringBootApplication
@EnableKafka
class Application {
@KafkaListener(topics = ["test"])
fun onMessage(message: String) {
log.warn("onMessage: $message")
}
拦截器:
class Interceptor : ConsumerInterceptor<String, String> {
override fun onCommit(offsets: MutableMap<TopicPartition, OffsetAndMetadata>) {
log.warn("onCommit: $offsets")
}
override fun onConsume(records: ConsumerRecords<String, String>): ConsumerRecords<String, String> {
log.warn("onConsume: $records")
return records
}
}
应用配置:
spring:
kafka:
consumer:
enable-auto-commit: false
auto-offset-reset: earliest
properties:
"interceptor.classes": com.example.Interceptor
group-id: test-group
listener:
ack-mode: record
在使用@EmbeddedKafka 的测试中:
@Test
fun sendMessage() {
kafkaTemplate.send("test", "id", "sent message").get() // block so we don't end before the consumer gets the message
}
这输出了我所期望的:
onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@6a646f3c
onMessage: sent message
onCommit: {test-0=OffsetAndMetadata{offset=1, leaderEpoch=null, metadata=''}}
但是,当我通过提供transaction-id-prefix 启用事务时,不再调用Interceptor 的onCommit。
我更新的配置仅添加:
spring:
kafka:
producer:
transaction-id-prefix: tx-id-
并且测试更新为将send 包装在事务中:
@Test
fun sendMessage() {
kafkaTemplate.executeInTransaction {
kafkaTemplate.send("test", "a", "sent message").get()
}
}
通过此更改,我的日志输出现在只有
onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@738b5968
onMessage: sent message
调用Interceptor 的onConsume 方法,@KafkaListener 接收消息,但从未调用onCommit。
有谁碰巧知道这里发生了什么?我对我应该在这里看到的内容的期望是否不正确?
【问题讨论】:
标签: spring-boot kotlin apache-kafka spring-kafka