package com.sea.cbb.utils;
import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.*;
/************************************
*
* org.apache.kafka
* kafka-clients
* 3.0.0
*
* @PACKAGE : com.sea.cbb.utils
* @Author : Sea
* @Date : 2022/6/26 17:23
* @Desc :
* @History :
***********************************/
public class KafkaConsumerUtils {
@FunctionalInterface
public interface KafkaBacthMsgHandler{
void handler(KafkaConsumer consumer,ConsumerRecords records) throws Exception;
}
@FunctionalInterface
public interface KafkaMsgHandler{
Out handler(ConsumerRecord record);
}
/**
* test
* @param args
*/
public static void main(String[] args) {
// KafkaConsumerUtils.consumerMsgAutoCommit("cbb_api_request_log","192.168.18.54:9092,192.168.18.199:9092,192.168.18.176:9092",
// "test",
// (record) -> {
// System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
// });
KafkaConsumerUtils. bacthMsgManualCommitSync("cbb_api_request_log", "192.168.18.54:9092,192.168.18.199:9092,192.168.18.176:9092",
"test11", (KafkaConsumer consumer, ConsumerRecords records) ->{
for(Object record1 : records){
ConsumerRecord record = (ConsumerRecord) record1;
System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
}
consumer.commitSync();
}
);
}
private static Properties initSerializer() {
Properties pros = new Properties();
pros.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");
pros.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer");
pros.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer");
return pros;
}
/**
* consumer 1 : 自动提交位移
*/
private static void consumerMsgAutoCommit(String topic,String bootstrapServers,String group,KafkaMsgHandler kafkaMagHandler){
Properties pros = initSerializer();
pros.put("bootstrap.servers",bootstrapServers);
pros.put("group.id",group);
pros.put("enable.auto.commit",true);
KafkaConsumer consumer = new KafkaConsumer(pros);
consumer.subscribe(Collections.singletonList(topic));
//指定topic
// TopicPartition partition = new TopicPartition(topic, 1);
// List lists = Arrays.asList(partition);
// consumer.assign(lists);
// consumer.seekToBeginning(lists);
//// consumer.seek(partition, 0);
try{
while(true){
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
Iterator iterator = records.iterator();
while (iterator.hasNext()){
ConsumerRecord record = (ConsumerRecord) iterator.next();
kafkaMagHandler.handler(record);
// System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
}
try{
Thread.sleep(1000);
}catch (Exception e){
e.printStackTrace();
}
}
}finally {
consumer.close();
}
}
/**
* consumer 2 : 手动提交位移, 批量提交
*/
public static void bacthMsgManualCommitSync(String topic,String bootstrapServers,String group,KafkaBacthMsgHandler kafkaMsgHandler) {
Properties pros = initSerializer();
pros.put("bootstrap.servers",bootstrapServers);
pros.put("group.id",group);
pros.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);//获取最大提交数量1
pros.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,false);
pros.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");
KafkaConsumer consumer = new KafkaConsumer(pros);
consumer.subscribe(Collections.singletonList(topic));
while(true)
{
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
try{
kafkaMsgHandler.handler(consumer,records);
// for(Object record1 : records){
// ConsumerRecord record = (ConsumerRecord) record1;
// System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
// }
// consumer.commitSync();
}catch (Exception e){
e.printStackTrace();
System.out.println("commit failed msg" + e.getMessage());
}
}
}
/**
* consumer 3 异步提交位移
*/
public static void consumerMsgManualCommitAsync(String topic,String bootstrapServers,String group){
Properties pros = initSerializer();
pros.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);//获取最大提交数量1
pros.put("bootstrap.servers",bootstrapServers);
pros.put("group.id",group);
pros.put("enable.auto.commit",false);
KafkaConsumer consumer = new KafkaConsumer<>(pros);
consumer.subscribe(Collections.singletonList(topic));
while(true)
{
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
for(Object record1 : records){
ConsumerRecord record = (ConsumerRecord) record1;
System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
}
consumer.commitAsync();
}
}
/**
* consumer 4 异步提交位移带回调
*/
public static void consumerMessageManualCommitAsyncWithCallBack(String topic,String bootstrapServers,String group){
Properties pros = initSerializer();
pros.put("bootstrap.servers",bootstrapServers);
pros.put("group.id",group);
pros.put("enable.auto.commit",false);
KafkaConsumer consumer = new KafkaConsumer<>(pros);
consumer.subscribe(Collections.singletonList("kafkatest"));
while(true){
ConsumerRecords records = consumer.poll(80);
for(Object record1 : records){
ConsumerRecord record = (ConsumerRecord) record1;
System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
}
consumer.commitAsync((offsets,e)->{
if(null != e){
System.out.println("commit async callback error" + e.getMessage());
System.out.println(offsets);
}
});
}
}
/**
* consumer 5 混合提交方式
*/
public static void mixSyncAndAsyncCommit(String topic,String bootstrapServers,String group){
Properties pros = initSerializer();
pros.put("bootstrap.servers",bootstrapServers);
pros.put("group.id",group);
pros.put("enable.auto.commit",false);
KafkaConsumer consumer = new KafkaConsumer<>(pros);
consumer.subscribe(Collections.singletonList(topic));
try{
while(true){
ConsumerRecords records = consumer.poll(Duration.ofMillis(100));
for(Object record1 : records){
ConsumerRecord record = (ConsumerRecord) record1;
System.out.println(record.timestamp() + "," +record.topic() + "," + record.partition() + "," + record.offset() + " " + record.key() +"," + record.value());
}
consumer.commitAsync();
}
}catch (Exception e){
System.out.println("commit async error: " + e.getMessage());
}finally {
try{
consumer.commitSync();
}finally {
consumer.close();
}
}
}
}