Spring-boot整合Kafka

生产者

说明

KafkaTemplate封装了一个生成器,并提供了方便的方法来发送数据到kafka主题。 提供了异步和同步方法,异步方法返回一个Future。

其构造方法有:

    ListenableFuture<SendResult<K, V>> sendDefault(V data);

    ListenableFuture<SendResult<K, V>> sendDefault(K key, V data);

    ListenableFuture<SendResult<K, V>> sendDefault(int partition, K key, V data);

    ListenableFuture<SendResult<K, V>> send(String topic, V data);

    ListenableFuture<SendResult<K, V>> send(String topic, K key, V data);

    ListenableFuture<SendResult<K, V>> send(String topic, int partition, V data);

    ListenableFuture<SendResult<K, V>> send(String topic, int partition, K key, V data);

    ListenableFuture<SendResult<K, V>> send(Message<?> message);

前3个方法需要向Temple提供默认主题

配置

使用Producer配置类

@Configuration
@EnableKafka
public class ProducerConfig {

    @Value("${kafka.producer.servers}")
    private String servers;

    @Value("${kafka.producer.retries}")
    private int retries;

    @Value("${kafka.producer.batch.size}")
    private int batchSize;

    @Value("${kafka.producer.linger}")
    private int linger;

    @Value("${kafka.producer.buffer.memory}")
    private int bufferMemory;

    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(Config.BOOTSTRAP_SERVERS_CONFIG, servers);
        props.put(Config.RETRIES_CONFIG, retries);
        props.put(Config.BATCH_SIZE_CONFIG, batchSize);
        props.put(Config.LINGER_MS_CONFIG, linger);
        props.put(Config.BUFFER_MEMORY_CONFIG, bufferMemory);
        props.put(Config.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(Config.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(Config.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return props;
    }

    public ProducerFactory<String, String> producerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<String, String>(producerFactory());
    }
}

示例

@RestController
@RequestMapping("/kafka/producer")
public class ProducerController {
    private static Logger logger = LoggerFactory.getLogger(ProducerController.class);

    @Value("${topic.name}")
    private String topicName;

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @RequestMapping("/send")
    public Object sendKafka(String message) {
        try {
            logger.info("send kafka message: {}", message);
            kafkaTemplate.send(topicName, UUID.randomUUID().toString(), message);
            return "success";
        } catch (Exception e) {
            logger.error("发送kafka失败", e);
            return "fail";
        }
    }
}

消费者

说明

可以通过配置MessageListenerContainer并提供MessageListener或通过使用@KafkaListener注释来接收消息。
MessageListenerContainer有两个实现:

  • KafkaMessageListenerContainer:从单个线程上的所有主题/分区接收所有消息
  • ConcurrentMessageListenerContainer:委托给1个或多个KafkaMessageListenerContainer以提供多线程消费。通过container.setConcurrency(3),来设置多个线程

配置

使用Consumer配置类

@Configuration
@EnableKafka
public class ConsumerConfig {

    @Value("${kafka.consumer.servers}")
    private String servers;

    @Value("${kafka.consumer.enable.auto.commit}")
    private boolean enableAutoCommit;

    @Value("${kafka.consumer.session.timeout}")
    private String sessionTimeout;

    @Value("${kafka.consumer.auto.commit.interval}")
    private String autoCommitInterval;

    @Value("${kafka.consumer.group.id}")
    private String groupId;

    @Value("${kafka.consumer.topic}")
    private String topic;

    @Value("${kafka.consumer.auto.offset.reset}")
    private String autoOffsetReset;

    @Value("${kafka.consumer.concurrency}")
    private int concurrency;

    /**
     * KafkaMessageListenerContainer: 从单个线程上的所有主题/分区接收所有消息

    @Bean(initMethod = "doStart")
    public KafkaMessageListenerContainer<String, String> kafkaMessageListenerContainer() {
        KafkaMessageListenerContainer<String, String> container = new KafkaMessageListenerContainer<>(consumerFactory(), containerProperties());
        return container;
    }

    */

    /**
     * ConcurrentMessageListenerContainer:
     * 委托给1个或多个KafkaMessageListenerContainer以提供多线程消费。
     * 通过container.setConcurrency(3),来设置多个线程
     */
    @Bean(initMethod = "doStart")
    public ConcurrentMessageListenerContainer<String, String> concurrentMessageListenerContainer() {
        ConcurrentMessageListenerContainer<String, String> container = new ConcurrentMessageListenerContainer<>(consumerFactory(), containerProperties());
        container.setConcurrency(concurrency);
        return container;

    }

    public ConsumerFactory<String, String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    public ContainerProperties containerProperties() {
        ContainerProperties containerProperties = new ContainerProperties(topic);
        containerProperties.setMessageListener(messageListener());
        return containerProperties;
    }

    public Map<String, Object> consumerConfigs() {
        Map<String, Object> propsMap = new HashMap<>();
        propsMap.put(Config.BOOTSTRAP_SERVERS_CONFIG, servers);
        propsMap.put(Config.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit);
        propsMap.put(Config.AUTO_COMMIT_INTERVAL_MS_CONFIG, autoCommitInterval);
        propsMap.put(Config.SESSION_TIMEOUT_MS_CONFIG, sessionTimeout);
        propsMap.put(Config.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(Config.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(Config.GROUP_ID_CONFIG, groupId);
        propsMap.put(Config.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
        return propsMap;
    }

    public MessageListener<String, String> messageListener() {
        return new CustomMessageListener();
    }
}

消息接收

Java实现

直接使用kafka0.10 client去收发消息

@Test
public void receive(){
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
    props.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
    props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
    try{
        consumer.subscribe(Arrays.asList(topic));
        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(10000);
            records.forEach(record -> {
                System.out.printf("client : %s , topic: %s , partition: %d , offset = %d, key = %s, value = %s%n", clientId, record.topic(),
                        record.partition(), record.offset(), record.key(), record.value());
            });
        }
    }catch (Exception e){
        e.printStackTrace();
    }finally {
        consumer.close();
    }
}
使用MessageListener接口

继承MessageListener接口

public class CustomMessageListener implements MessageListener<Integer, String> {
    private static Logger logger = LoggerFactory.getLogger(CustomMessageListener.class);

    @Override
    public void onMessage(ConsumerRecord<Integer, String> data) {
        logger.info("received key: {}, value: {}", data.key(), data.value());
    }

  //或包含消费者的onMessage方法,以手动提交ofset
}
使用@KafkaListener注解
@KafkaListener(id = "foo", topics = "myTopic")
public void listen(String data) {
     ...
}

@KafkaListener(id = "bar", topicPartitions =
        { @TopicPartition(topic = "topic1", partitions = { "0", "1" }),
          @TopicPartition(topic = "topic2", partitions = "0",
             partitionOffsets = @PartitionOffset(partition = "1", initialOffset = "100"))
        })
public void listen(ConsumerRecord<?, ?> record) {
    ...
}

总结

  • 对于生产者来说,封装KafkaProducer到KafkaTemplate相对简单
  • 对于消费者来说,由于spring是采用注解的形式去标注消息处理方法
    1. 先在KafkaListenerAnnotationBeanPostProcessor中扫描bean,然后注册到KafkaListenerEndpointRegistrar
    2. 而KafkaListenerEndpointRegistrar在afterPropertiesSet的时候去创建MessageListenerContainer
    3. messageListener包含了原始endpoint携带的bean以及method转换成的InvocableHandlerMethod
    4. ConcurrentMessageListenerContainer这个衔接上,根据配置的spring.kafka.listener.concurrency来生成多个并发的KafkaMessageListenerContainer实例
    5. 每个KafkaMessageListenerContainer都自己创建一个ListenerConsumer,然后自己创建一个独立的kafka consumer,每个ListenerConsumer在线程池里头运行,这样来实现并发
    6. 每个ListenerConsumer里头都有一个recordsToProcess队列,从原始的kafka consumer poll出来的记录会放到这个队列里头,
    7. 然后有一个ListenerInvoker线程循环超时等待从recordsToProcess取出记录,然后调用messageListener的onMessage方法(即KafkaListener注解标准的方法)
项目源码

https://github.com/scjqwe/spring-kafka-examples

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

推荐阅读更多精彩内容