【发布时间】:2019-06-12 21:13:21
【问题描述】:
有时(似乎很随机)Kafka 会发送旧消息。我只想要最新的消息,所以它会用相同的键覆盖消息。目前看起来我有多条消息具有相同的密钥,它没有被压缩。
我在主题中使用这个设置:
cleanup.policy=compact
我正在使用 Java/Kotlin 和 Apache Kafka 1.1.1 客户端。
Properties(8).apply {
val jaasTemplate = "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"%s\" password=\"%s\";"
val jaasCfg = String.format(jaasTemplate, Configuration.kafkaUsername, Configuration.kafkaPassword)
put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
BOOTSTRAP_SERVERS)
put(ConsumerConfig.GROUP_ID_CONFIG,
"ApiKafkaKotlinConsumer${Configuration.kafkaGroupId}")
put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer::class.java.name)
put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer::class.java.name)
put("security.protocol", "SASL_SSL")
put("sasl.mechanism", "SCRAM-SHA-256")
put("sasl.jaas.config", jaasCfg)
put("max.poll.records", 100)
put("receive.buffer.bytes", 1000000)
}
我错过了一些设置吗?
【问题讨论】:
-
即使您设置了下面答案中提到的所有属性,压缩仍然作为后台进程运行以删除重复键。所以消费者仍然会看到一些或所有重复的键。 Kafka 不会像您想象的那样即时过滤重复的键。你可以使用 Kafka Streams 来实现你想要的——this 可以帮助你。
标签: java kotlin apache-kafka kafka-consumer-api