【问题标题】:Kafks consumer.poll returns no dataKafka consumer.poll 不返回任何数据
【发布时间】:2019-07-26 01:09:57
【问题描述】:

我有两个 Kafka (2.11-0.11.0.1) 代理。主题的默认复制因子设置为 2。生产者仅将数据写入零分区。

我已经安排了定期运行任务的执行程序。当它每分钟使用少量记录(每分钟 100 条)的主题时,就会像魅力一样工作。但是对于大型主题(每分钟 10K)方法 poll 不返回任何数据。

任务是:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public final class TopicToDbPump implements Runnable {
  private static final Logger log = LoggerFactory.getLogger(TopicToDbPump.class);
  private final String topic;
  private final TopicPartition topicPartition;
  private final Properties properties;

  public TopicToDbPump(String topic, Properties properties) {
    this.topic = topic;
    topicPartition = new TopicPartition(topic, 0);
    this.properties = properties;
  }

  @Override
  public void run() {
    try (final Consumer<String, String> consumer = new KafkaConsumer<>(properties)) {
      consumer.assign(Collections.singleton(topicPartition));
      final long offset = readOffsetFromDb(topic);
      consumer.seek(topicPartition, offset);
      final ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
      if (records.isEmpty()) {
        log.debug("No data from topic " + topic + " available");
        return;
      }
      saveData(records.records(topic));
    } catch (Throwable t) {
      log.error("Etl process " + topic + " failed with exception", t);
    }
  }
}

消费者参数为:

"bootstrap.servers" = "host-1:9092,host-2:9092",
"group.id" = "my-group",
"enable.auto.commit" = "false",
"key.deserializer" = "org.apache.kafka.common.serialization.StringDeserializer",
"value.deserializer" = "org.apache.kafka.common.serialization.StringDeserializer",
"max.partition.fetch.bytes": "50000000",
"max.poll.records": "10000"

怎么了?

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    Kafka Consumer API 不保证第一次调用poll() 会返回任何数据。

    消费者首先必须连接到集群,发现分配给它的所有分区的领导者。正如您想象的那样,这可能需要几秒钟,因此数据不太可能立即到达。

    如果没有首先返回数据,您应该多次调用poll()。

    【讨论】:

    • 感谢您的解释!我知道在分配分区和使用subscribe 方法在poll 中获得结果之间存在超时。但不知何故,我认为手动分配分区会导致立即获取数据。
    猜你喜欢
    • 2019-09-20
    • 1970-01-01
    • 2018-10-19
    • 2017-11-23
    • 2021-12-13
    • 2016-12-06
    • 2011-04-04
    • 2015-06-09
    • 1970-01-01
    相关资源
    最近更新 更多