kafkaUtils (原生java API consumer)


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();

            }

        }

    }

}