springboot集成Kafka kafka和MQ的区别:
本文内容纲要:
-kafka和MQ的区别:
-Springboot集成kafka代码展示:
kafka和MQ的区别:
1)在架构模型方面,
RabbitMQ遵循AMQP协议,RabbitMQ的broker由Exchange,Binding,queue组成,其中exchange和binding组成了消息的路由键;客户端Producer通过连接channel和server进行通信,Consumer从queue获取消息进行消费(长连接,queue有消息会推送到consumer端,consumer循环从输入流读取数据)。rabbitMQ以broker为中心;有消息的确认机制。
kafka遵从一般的MQ结构,producer,broker,consumer,以consumer为中心,消息的消费信息保存的客户端consumer上,consumer根据消费的点,从broker上批量pull数据;无消息确认机制。
2)在吞吐量,
rabbitMQ在吞吐量方面稍逊于kafka,他们的出发点不一样,rabbitMQ支持对消息的可靠的传递,支持事务,不支持批量的操作;基于存储的可靠性的要求存储可以采用内存或者硬盘。
kafka具有高的吞吐量,内部采用消息的批量处理,zero-copy机制,数据的存储和获取是本地磁盘顺序批量操作,具有O(1)的复杂度,消息处理的效率很高。
3)在可用性方面,
rabbitMQ支持miror的queue,主queue失效,mirorqueue接管。
kafka的broker支持主备模式。
4)在集群负载均衡方面,
rabbitMQ的负载均衡需要单独的loadbalancer进行支持。
kafka采用zookeeper对集群中的broker、consumer进行管理,可以注册topic到zookeeper上;通过zookeeper的协调机制,producer保存对应topic的broker信息,可以随机或者轮询发送到broker上;并且producer可以基于语义指定分片,消息发送到broker的某分片上。
转载:https://www.cnblogs.com/csuliujia/p/9379402.html
Springboot集成kafka代码展示:
对象:
packagecom.example.demo;
importcom.alibaba.fastjson.JSON;
importorg.springframework.beans.factory.annotation.Autowired;
importorg.springframework.kafka.core.KafkaTemplate;
importorg.springframework.stereotype.Component;
@Component
publicclassProduct{
@Autowired
privateKafkaTemplatekafkaTemplate;
publicvoidsend(Stringname){
Useru=newUser();
u.setName(name);
u.setAge(11);
kafkaTemplate.send("user",JSON.toJSONString(u));
}
}
生产者:
packagecom.example.demo;
importcom.alibaba.fastjson.JSON;
importorg.springframework.beans.factory.annotation.Autowired;
importorg.springframework.kafka.core.KafkaTemplate;
importorg.springframework.stereotype.Component;
@Component
publicclassProduct{
@Autowired
privateKafkaTemplatekafkaTemplate;
publicvoidsend(Stringname){
Useru=newUser();
u.setName(name);
u.setAge(11);
kafkaTemplate.send("user",JSON.toJSONString(u));
}
}
消费者:
packagecom.example.demo;
importorg.apache.kafka.clients.consumer.ConsumerRecord;
importorg.springframework.kafka.annotation.KafkaListener;
importorg.springframework.stereotype.Component;
importjava.util.Optional;
@Component
publicclassConsumer{
@KafkaListener(topics="user")
publicvoidconsumer(ConsumerRecordconsumerRecord){
Optional<Object>kafkaMassage=Optional.ofNullable(consumerRecord.value());
if(kafkaMassage.isPresent()){
Objecto=kafkaMassage.get();
System.out.println(o);
}
}
}
测试:
packagecom.example.demo;
importorg.springframework.beans.factory.annotation.Autowired;
importorg.springframework.boot.SpringApplication;
importorg.springframework.boot.autoconfigure.SpringBootApplication;
importjavax.annotation.PostConstruct;
@SpringBootApplication
publicclassDemoApplication{
@Autowired
privateProductproduct;
@PostConstruct
publicvoidinit(){
for(inti=0;i<10;i++){
product.send("afs"+i);
}
}
publicstaticvoidmain(String[]args){
SpringApplication.run(DemoApplication.class,args);
}
}
配置文件:application.properties
spring.application.name=kafka-user
server.port=8080
#==============kafka===================
#指定kafka代理地址,可以多个
spring.kafka.bootstrap-servers=localhost:9092
#===============provider=======================
spring.kafka.producer.retries=0
#每次批量发送消息的数量
spring.kafka.producer.batch-size=16384
spring.kafka.producer.buffer-memory=33554432
#指定消息key和消息体的编解码方式
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
#===============consumer=======================
#指定默认消费者groupid
spring.kafka.consumer.group-id=user-log-group
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.enable-auto-commit=true
spring.kafka.consumer.auto-commit-interval=100
#指定消息key和消息体的编解码方式
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
配置文件:pom.xml
<?xmlversion="1.0"encoding="UTF-8"?>
<projectxmlns="http://maven.apache.org/POM/4.0.0"xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.1.6.RELEASE</version>
<relativePath/><!--lookupparentfromrepository-->
</parent>
<groupId>com.example</groupId>
<artifactId>demo</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>demo</name>
<description>DemoprojectforSpringBoot</description>
<properties>
<java.version>1.8</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!--引入kafak和spring整合的jar-->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>2.2.7.RELEASE</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.44</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
注意:此处有版本问题,springboot2.x得用kafka2.x
另外:记得开启zookeeper+kafka
安装参考:https://www.cnblogs.com/lnice/p/9668750.html(windows版本)
https://blog.csdn.net/xuzhelin/article/details/71515208(Linux版本)
本文内容总结:kafka和MQ的区别:,Springboot集成kafka代码展示:,
原文链接:https://www.cnblogs.com/tysl/p/11170811.html