kafka消费者工具类

手写编写一个工具类的接口程序:
package com.ky.kafka.consumer;

/**

  • @Author: xwj

  • @Date: 2019/3/12 0012 09:48

  • @Version 1.0
    */
    public interface ConsumerInterface {

    /**

    • 普通的消费者
    • @param topic 主题
      */
      public void consumer(String topic);

    /**

    • 批量消费
    • @param topic 主题
    • @param num 批量消费的数量
      */
      public void consumerBatch(String topic, int num);

    /**

    • 指定从哪里开始消费数据
    • @param topic 主题
    • @param line 具体的位置
      */
      public void consumerFromDetailPostion(String topic, long line);
      }

其次编写一个工厂类生产实例:
import java.util.Properties;

/**

  • @Author: xwj

  • @Date: 2019/3/12 0012 09:56

  • @Version 1.0
    */
    public class ConsumerFactory {

    /**

    • 获取消费者实例
    • @return 获取消费者实例
      */
      public static KafkaConsumer getConsumer() {
      Properties props = new Properties();
      props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, PropertyUtil.getInstance().getString("brokerList", ""));
      props.put(ConsumerConfig.GROUP_ID_CONFIG, PropertyUtil.getInstance().getString("groupid", "123"));
      props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
      props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");
      props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "30000");
      props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
      props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
      props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
      System.out.println("初始化kafka消费者的配置信息......");
      return new KafkaConsumer(props);
      }
      }

最后编写一个工具类通用的访问方法:
package com.ky.kafka.consumer;

import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Collection;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

/**

  • @Author: xwj

  • @Date: 2019/3/12 0012 09:59

  • @Version 1.0
    */
    public class Consumer implements ConsumerInterface {
    private final static Logger logger = LoggerFactory.getLogger(Consumer.class);

    @Override
    public void consumer(String topic) {
    KafkaConsumer consumer = ConsumerFactory.getConsumer();
    consumer.subscribe(Collections.singletonList(topic));
    try {
    while (true) {
    ConsumerRecords<String, String> records = consumer.poll(300);
    for (ConsumerRecord<String, String> record : records) {
    System.out.println(String.format("partition = %d , offset = %d, key = %s, value = %s", record.partition(),
    record.offset(), record.key(), record.value()));
    }
    consumer.commitAsync(new OffsetCommitCallback() {
    @Override
    public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e) {
    if (e != null) {
    logger.error("commit failed for offset {}", offsets, e);
    }
    }
    });
    }
    } finally {
    try {
    consumer.commitSync();
    } finally {
    release(consumer);
    }
    }
    }

    @Override
    public void consumerBatch(String topic, int num) {
    KafkaConsumer consumer = ConsumerFactory.getConsumer();
    ConcurrentHashMap<TopicPartition, OffsetAndMetadata> currentOffsets = new ConcurrentHashMap<TopicPartition, OffsetAndMetadata>();
    int count = 0;
    consumer.subscribe(Collections.singletonList(topic));
    try {
    while (true) {
    final ConsumerRecords<String, String> consumerRecords = consumer.poll(300);
    for (ConsumerRecord<String, String> consumerRecord : consumerRecords) {
    System.out.println(consumerRecord.value());
    currentOffsets.put(new TopicPartition(consumerRecord.topic(), consumerRecord.partition()), new OffsetAndMetadata(consumerRecord.offset() + 1, "no metadata"));
    if (count % num == 0) {
    consumer.commitAsync(currentOffsets, new OffsetCommitCallback() {
    @Override
    public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e) {
    if (e != null) {
    logger.error("commit failed for offset {}", offsets, e);
    }
    }
    });
    }
    count++;
    }
    }
    } catch (Exception e) {
    e.printStackTrace();
    } finally {
    try {
    consumer.commitSync();
    } finally {
    release(consumer);
    }
    }
    }

    @Override
    public void consumerFromDetailPostion(String topic, final long line) {
    final KafkaConsumer consumer = ConsumerFactory.getConsumer();
    final ConcurrentHashMap<TopicPartition, OffsetAndMetadata> currentOffsets = new ConcurrentHashMap<TopicPartition, OffsetAndMetadata>();
    try {
    consumer.subscribe(Collections.singletonList(topic), new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> collection) {
    consumer.commitSync(currentOffsets);
    }

             @Override
             public void onPartitionsAssigned(Collection<TopicPartition> collection) {
                 //consumer.seekToBeginning(collection);
                 for (TopicPartition topicPartition : collection) {
                     consumer.seek(topicPartition, line);
                 }
    
             }
         });
    
         while (true) {
             final ConsumerRecords<String, String> records = consumer.poll(300);
             for (ConsumerRecord<String, String> record : records) {
                 System.out.println(String.format("partition = %d , offset = %d, key = %s, value = %s", record.partition(),
                         record.offset(), record.key(), record.value()));
                 currentOffsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1, "no metadata"));
             }
             consumer.commitAsync(currentOffsets, new OffsetCommitCallback() {
                 @Override
                 public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e) {
                     if (e != null) {
                         logger.error("commit failed for offset {}", offsets, e);
                     }
                 }
             });
         }
     } catch (Exception e) {
         e.printStackTrace();
     } finally {
         try {
             consumer.commitSync();
         } finally {
             release(consumer);
         }
     }
    

    }

    /**

    • 关闭消费者

    • @param consumer 消费者
      */
      private void release(final KafkaConsumer consumer) {
      Runtime.getRuntime().addShutdownHook(new Thread() {

       @Override
       public void run() {
           super.run();
           if (consumer != null) {
               consumer.close();
           }
       }
      

      });
      }
      }

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 199,340评论 5 467
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 83,762评论 2 376
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 146,329评论 0 329
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 53,678评论 1 270
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 62,583评论 5 359
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 47,995评论 1 275
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 37,493评论 3 390
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,145评论 0 254
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 40,293评论 1 294
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,250评论 2 317
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,267评论 1 328
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 32,973评论 3 316
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 38,556评论 3 303
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,648评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 30,873评论 1 255
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 42,257评论 2 345
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 41,809评论 2 339

推荐阅读更多精彩内容