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