天天看點

Kafka2.0消費者用戶端使用1 初始化配置2 訂閱主題3 拉取消息4 送出偏移量5 其他方法

1 初始化配置

  Kafka 通過 KafkaConsumer 構造器初始化生産者用戶端的配置。

  常用的重要配置,詳見官網。

  • bootstrap.servers:Kafka 叢集位址(host1:post,host2:post),Kafka 用戶端初始化時會自動發現位址,是以可以不填寫所有位址。
  • group.id:消費組 ID
  • key.serializer:實作了 Kafka 序列化接口的類,用來序列化 key。
  • value.serializer:實作了 Kafka 序列化接口的類,用來序列化 value。
  • enable.auto.commit:預設 true,表示消費者偏移量會定期送出到背景。
  • auto.offset.reset:Kafka 的偏移量。

     earliest:自動重置為最早的偏移量。

     latest:自動重置為最新的偏移量。

     none:如果沒有找到消費組之前的那個偏移量,則消費者抛出異常。

     其他:消費者抛出異常。

  • fetch.min.bytes/fetch.max.bytes:消費者一次拉取的最小/最大值。
  • max.poll.interval.ms:消費者拉取的最大間隔時間,逾時後從組中移除消費者。
  • heartbeat.interval.ms:心跳發送間隔的逾時時間,逾時後從組中移除消費者。
  • isolation.level:事務的隔離級别。

     read_uncommitted:預設,可以消費到所有消息,包括被中止的消息。

     read_committed:隻能消費到事務送出過的消息。

     非事務性消息無條件傳回。

// 基礎配置
Map<String, Object> configs = new HashMap<>();
configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092,127.0.0.1:9093,127.0.0.1:9094");
configs.put(ConsumerConfig.GROUP_ID_CONFIG, "my_test");
configs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(configs);
           

2 訂閱主題

  Kafka 消費者提供4種方式訂閱主題,1種方式指定分區。

  • topics:指定主題集。
  • pattern:指定正規表達式來比對主題。
  • listener:消費者再均衡監聽器。
  • partitions:指定分區集合。
// 指定主題
public void subscribe(Collection<String> topics, ConsumerRebalanceListener listener)
public void subscribe(Collection<String> topics)
public void subscribe(Pattern pattern, ConsumerRebalanceListener listener)
public void subscribe(Pattern pattern)
// 指定分區
public void assign(Collection<TopicPartition> partitions)
           

3 拉取消息

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(configs);
consumer.subscribe(Collections.singletonList("test")); // 指定主題
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(3));
           

 3.1 指定偏移量消費

TopicPartition tp = new TopicPartition("test", 0);
consumer.assign(Collections.singletonList(tp)); // 訂閱指定分區
consumer.seek(tp, 4L); // 指定分區偏移量值為4
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(3));
           

 3.2 指定時間戳消費

TopicPartition tp = new TopicPartition("test", 0);
        consumer.assign(Collections.singletonList(tp)); // 訂閱指定分區
        Map<TopicPartition, Long> tpTime = new HashMap<>();
        tpTime.put(tp, 1563027475113L); // 指定時間戳
        Map<TopicPartition, OffsetAndTimestamp> tpOffsetAndTime = consumer.offsetsForTimes(tpTime);
        long offset = tpOffsetAndTime.get(tp).offset(); // 擷取偏移量
        consumer.seek(tp, offset); // 指定偏移量
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(3));
           

4 送出偏移量

 4.1 同步送出

  參數說明

  • offsets:可以指定送出分區的偏移量。
  • timeout:偏移量送出成功的逾時時間。
public void commitSync()
public void commitSync(Duration timeout)
public void commitSync(final Map<TopicPartition, OffsetAndMetadata> offsets)
public void commitSync(final Map<TopicPartition, OffsetAndMetadata> offsets, final Duration timeout)
           

 4.2 異步送出

  參數說明

  • offsets:可以指定送出分區的偏移量。
  • callback:異步回調。
public void commitAsync()
public void commitAsync(OffsetCommitCallback callback)
public void commitAsync(final Map<TopicPartition, OffsetAndMetadata> offsets, OffsetCommitCallback callback)
           

5 其他方法

// 擷取配置設定給目前消費者的分區集合
public Set<TopicPartition> assignment()
// 取消訂閱
public void unsubscribe()
// 找到指定分區的第一個偏移量
public void seekToBeginning(Collection<TopicPartition> partitions)
// 找到指定分區的最後一個偏移量
public void seekToEnd(Collection<TopicPartition> partitions)
// 擷取指定分區即将消費的下一個偏移量
public long position(TopicPartition partition)
// 擷取指定分區最後送出的偏移量
public OffsetAndMetadata committed(TopicPartition partition)
// 擷取指定主題的分區清單
public List<PartitionInfo> partitionsFor(String topic)
// 擷取所有主題的資訊
public Map<String, List<PartitionInfo>> listTopics()
// 暫停消費
public void pause(Collection<TopicPartition> partitions)
// 恢複被暫停的消費
public void resume(Collection<TopicPartition> partitions)
// 擷取暫停的分區清單
public Set<TopicPartition> paused()
// 擷取指定分區第一個偏移量
public Map<TopicPartition, Long> beginningOffsets(Collection<TopicPartition> partitions)
// 擷取指定分區最後一個偏移量
public Map<TopicPartition, Long> endOffsets(Collection<TopicPartition> partitions)
// 喚醒消費者
public void wakeup()