Kafka+SpringBoot整合

1.SpringBoot项目导入Kafka依赖(版本问题需要注意下springboot1.5以后的可以对应kafka2.2)

<dependency>
          <groupId>org.springframework.kafka</groupId>
          <artifactId>spring-kafka</artifactId>
          <version>2.2.5.RELEASE</version>
</dependency>

2.application.yml 配置kafka

spring:
  kafka:
    bootstrap-servers: localhost:9092
    consumer:
      group-id: 02
      auto-offset-reset: earliest
      enable-auto-commit: false
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
    producer:
      acks: 1 #收到确认数
      batch-size: 16384
      buffer-memory: 33554432
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      retries: 1 #重发次数,设置大于0,则失败重发

3.生产者向Kafka发送消息

@Autowired
    private KafkaTemplate kafkaTemplate;

    public void sendMessage(String orderId){
        Order order = new Order();
        order.setOrderId(orderId).setOrderContent("下单了").setCreateTime(new Date());
        kafkaTemplate.send("order-dispacher3",0,"1",JSON.toJSONString(order));
        kafkaTemplate.send("order-dispacher3", 1,"2",JSON.toJSONString(order));
        //kafkaTemplate.send("order-dispacher", JSON.toJSONString(order));
        logger.info("kafka消息发送成功"+order.toString());
    }

使用kafkaTemplate发送消息,send函数重载了多个方法,如下图,可以指定partition(分区),key发送消息。


image.png

4.消费者向Kafka消费消息

@KafkaListener(topics = "order-dispacher3")
    public void readMessage(ConsumerRecord<String,String> consumerRecord){
        logger.info("消费消息:"+consumerRecord.toString());
    }

调试中的爬坑记录:
问题:kafka Java创建生产者报错:Invalid partition given with record: 1 is not in the range [0...1)]
直接使用kafkaTemplete发送消息,是无法创建partition(分区),得提前创建好分区数,默认一个分区,或者修改分区数。partitions在是在创建topic的时候默认创建的partitions节点的个数,只对新创建的topic生效,所有尽量在项目规划时候定一个合理的值。也可以通过以下命令行动态扩容。

./bin/kafka-topics.bat --zookeeper  localhost:2181 --alter --partitions 2 --topic  test

5.附录Windows环境操作Kafka
5.1启动 (打开cmd,进入kafka文件目录)

cd D:\soft\kafka_2.12-2.1.0(kafka目录)
 .\bin\windows\kafka-server-start.bat .\config\server.properties

5.2创建topic

cd D:\soft\kafka_2.12-2.1.0\bin\windows
kafka-topics.bat --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic test

5.3打开producer

cd D:\soft\kafka_2.12-2.1.0\bin\windows
kafka-console-producer.bat --broker-list localhost:9092 --topic test

5.4打开consumer

cd D:\soft\kafka_2.12-2.1.0\bin\windows
kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic test --from-beginning
最后编辑于 :
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容