【问题标题】:Kafka doesn't save offset if consume short time如果消耗时间短,Kafka 不会保存偏移量
【发布时间】:2018-11-23 19:29:31
【问题描述】:

问题

具有特定组 id 的消费者连接到代理,监听主题不到 1 分钟并断开连接(根据业务逻辑)。当它收听主题时,它可以使用一些消息。 当同一个消费者重复此操作时,它会消费相同的消息!

我发现 Kafka 以 1 分钟的间隔保存偏移量。这意味着消费者必须听主题超过 1 分钟。 如何缩短此间隔?

我找到了这样的属性:

  • log.flush.offset.checkpoint.interval.ms
  • log.flush.start.offset.checkpoint.interval.ms
  • offset.flush.interval.ms - 看起来最合适

我尝试将它们设置在server.properties 文件中:

log.flush.offset.checkpoint.interval.ms=6000
log.flush.start.offset.checkpoint.interval.ms=6000
offset.flush.interval.ms=6000

重启 Kafka 和 Zookeeper。但这无济于事。消费者仍然需要听主题超过 1 分钟。我做错了什么?

我的环境

  • Kafka 和 Zookeeper 通过 Confluent。
  • php-rdkafka 作为客户端库
  • enable.auto.commit 设置为 true

我使用低级消费者。 auto.offset.reset 设置为 smallest。 代码示例

<?php
$topicConf = new \RdKafka\TopicConf();
$topicConf->set('auto.offset.reset', 'smallest');

$conf = new \RdKafka\Conf();
$conf->set('group.id', 'foo');

$kafkaConsumer = new \RdKafka\Consumer($conf);
$kafkaConsumer->addBrokers('queue.a:9092');
$kafkaConsumer->setLogLevel(LOG_DEBUG);

$topicConf = new \RdKafka\TopicConf();
$topicConf->set('auto.offset.reset', 'smallest');

$queue = $kafkaConsumer->newQueue();
$topic = $kafkaConsumer->newTopic('topic_name', $topicConf);
$topic->consumeQueueStart(0, \RD_KAFKA_OFFSET_STORED, $queue);

while (true) {
    $msg = $queue->consume(2000);
    if ($msg !== null) {
        var_dump($msg);
    }
}

【问题讨论】:

    标签: php apache-kafka


    【解决方案1】:

    您应该尝试在您的消费者中明确提交偏移量:

    在消费者中明确承诺抵消 如果您使用自动偏移提交,则无需担心显式提交偏移。但是,如果您决定需要更多地控制偏移提交的时间,您确实需要考虑如何提交偏移——要么是为了最大限度地减少重复,要么是因为您在主消费者轮询循环之外进行事件处理。

    摘自Kafka definitive guide,第 127 页。(这是一本免费的电子书,您可以下载)

    建议您始终在处理完事件后提交偏移量放轻松。您可以在轮询循环结束时使用自动提交配置或提交事件。

    我自己没用过php客户端,不过貌似this could be what you need

    添加到上面的代码示例:

    while (true) {
        $msg = $queue->consume(2000);
        if ($msg !== null) {
            var_dump($msg);
            $kafkaConsumer->commit($msg);
        }
    }
    

    【讨论】:

    • php-rdkafka 是一个有点奇怪的库。它仅对高级消费者类具有显式提交。我必须使用低级消费类。并且enable.auto.commit 设置为true。附:谢谢你的书:)
    • @EvgeniiKaravskii php-rdkafka 只是 librdkafka 库的一个包装器,它实际上完成了所有“繁重的工作”。它也用于其他语言。
    猜你喜欢
    • 1970-01-01
    • 2021-02-22
    • 2016-03-05
    • 2022-01-06
    • 1970-01-01
    • 2019-01-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多