【问题标题】:Kafka only subscribe to latest message?卡夫卡只订阅最新消息?
【发布时间】: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


【解决方案1】:

如果您希望每个键只有一个值,则必须使用 KTable<K,V> 抽象:StreamsBuilder::table(final String topic) from Kafka Streams。此处使用的主题应将清理策略设置为compact

如果您使用 KafkaConsumer,您只需从经纪人那里提取数据。它没有为您提供任何执行某种重复数据删除的机制。根据是否执行了压缩,您可以获得 onen 个相同键的消息。

关于压缩

压缩并不意味着同一键的所有旧值都会立即删除。当同一键的old 消息将被删除时,取决于几个属性。最重要的是:

  • log.cleaner.min.cleanable.ratio

符合清理条件的日志的脏日志与总日志的最小比率

  • log.cleaner.min.compaction.lag.ms

消息在日志中保持未压缩的最短时间。仅适用于正在压缩的日志。

  • log.cleaner.enable

启用日志清理程序以在服务器上运行。如果使用带有 cleanup.policy=compact 的任何主题,包括内部偏移主题,则应启用。如果禁用,这些主题将不会被压缩并会不断增长。

更多关于压缩的细节你可以找到https://kafka.apache.org/documentation/#compaction

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-10-12
    • 2017-06-12
    • 1970-01-01
    • 1970-01-01
    • 2021-01-26
    • 1970-01-01
    • 2021-04-02
    • 2016-11-03
    相关资源
    最近更新 更多