1 初始化配置
??Kafka 通過(guò) KafkaConsumer 構(gòu)造器初始化生產(chǎn)者客戶端的配置园爷。
??常用的重要配置戴质,詳見(jiàn)官網(wǎng)。
- bootstrap.servers:Kafka 集群地址(host1:post,host2:post)踢匣,Kafka 客戶端初始化時(shí)會(huì)自動(dòng)發(fā)現(xiàn)地址告匠,所以可以不填寫(xiě)所有地址。
- group.id:消費(fèi)組 ID
- key.serializer:實(shí)現(xiàn)了 Kafka 序列化接口的類(lèi)离唬,用來(lái)序列化 key后专。
- value.serializer:實(shí)現(xiàn)了 Kafka 序列化接口的類(lèi),用來(lái)序列化 value输莺。
- enable.auto.commit:默認(rèn) true戚哎,表示消費(fèi)者偏移量會(huì)定期提交到后臺(tái)裸诽。
- auto.offset.reset:Kafka 的偏移量。
?earliest:自動(dòng)重置為最早的偏移量型凳。
?latest:自動(dòng)重置為最新的偏移量丈冬。
?none:如果沒(méi)有找到消費(fèi)組之前的那個(gè)偏移量,則消費(fèi)者拋出異常甘畅。
?其他:消費(fèi)者拋出異常埂蕊。 - fetch.min.bytes/fetch.max.bytes:消費(fèi)者一次拉取的最小/最大值。
- max.poll.interval.ms:消費(fèi)者拉取的最大間隔時(shí)間疏唾,超時(shí)后從組中移除消費(fèi)者蓄氧。
- heartbeat.interval.ms:心跳發(fā)送間隔的超時(shí)時(shí)間,超時(shí)后從組中移除消費(fèi)者槐脏。
- isolation.level:事務(wù)的隔離級(jí)別喉童。
?read_uncommitted:默認(rèn),可以消費(fèi)到所有消息顿天,包括被中止的消息堂氯。
?read_committed:只能消費(fèi)到事務(wù)提交過(guò)的消息。
?非事務(wù)性消息無(wú)條件返回露氮。
// 基礎(chǔ)配置
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 消費(fèi)者提供4種方式訂閱主題祖灰,1種方式指定分區(qū)。
- topics:指定主題集畔规。
- pattern:指定正則表達(dá)式來(lái)匹配主題局扶。
- listener:消費(fèi)者再均衡監(jiān)聽(tīng)器。
- partitions:指定分區(qū)集合叁扫。
// 指定主題
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)
// 指定分區(qū)
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 指定偏移量消費(fèi)
TopicPartition tp = new TopicPartition("test", 0);
consumer.assign(Collections.singletonList(tp)); // 訂閱指定分區(qū)
consumer.seek(tp, 4L); // 指定分區(qū)偏移量值為4
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(3));
?3.2 指定時(shí)間戳消費(fèi)
TopicPartition tp = new TopicPartition("test", 0);
consumer.assign(Collections.singletonList(tp)); // 訂閱指定分區(qū)
Map<TopicPartition, Long> tpTime = new HashMap<>();
tpTime.put(tp, 1563027475113L); // 指定時(shí)間戳
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 同步提交
??參數(shù)說(shuō)明
- offsets:可以指定提交分區(qū)的偏移量三妈。
- timeout:偏移量提交成功的超時(shí)時(shí)間。
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 異步提交
??參數(shù)說(shuō)明
- offsets:可以指定提交分區(qū)的偏移量莫绣。
- callback:異步回調(diào)畴蒲。
public void commitAsync()
public void commitAsync(OffsetCommitCallback callback)
public void commitAsync(final Map<TopicPartition, OffsetAndMetadata> offsets, OffsetCommitCallback callback)
5 其他方法
// 獲取分配給當(dāng)前消費(fèi)者的分區(qū)集合
public Set<TopicPartition> assignment()
// 取消訂閱
public void unsubscribe()
// 找到指定分區(qū)的第一個(gè)偏移量
public void seekToBeginning(Collection<TopicPartition> partitions)
// 找到指定分區(qū)的最后一個(gè)偏移量
public void seekToEnd(Collection<TopicPartition> partitions)
// 獲取指定分區(qū)即將消費(fèi)的下一個(gè)偏移量
public long position(TopicPartition partition)
// 獲取指定分區(qū)最后提交的偏移量
public OffsetAndMetadata committed(TopicPartition partition)
// 獲取指定主題的分區(qū)列表
public List<PartitionInfo> partitionsFor(String topic)
// 獲取所有主題的信息
public Map<String, List<PartitionInfo>> listTopics()
// 暫停消費(fèi)
public void pause(Collection<TopicPartition> partitions)
// 恢復(fù)被暫停的消費(fèi)
public void resume(Collection<TopicPartition> partitions)
// 獲取暫停的分區(qū)列表
public Set<TopicPartition> paused()
// 獲取指定分區(qū)第一個(gè)偏移量
public Map<TopicPartition, Long> beginningOffsets(Collection<TopicPartition> partitions)
// 獲取指定分區(qū)最后一個(gè)偏移量
public Map<TopicPartition, Long> endOffsets(Collection<TopicPartition> partitions)
// 喚醒消費(fèi)者
public void wakeup()