H5W3
当前位置:H5W3 > java > 正文

【Java】springboot2.x整合kafka

springboot2.x整合kafka

一片秋叶一树春发布于 今天 10:38

pom引入

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

yml配置说明

spring:
kafka:
# 以逗号分隔的地址列表,用于建立与Kafka集群的初始连接(kafka 默认的端口号为9092)
bootstrap-servers: 192.168.1.200:9092
producer:
# 发生错误后,消息重发的次数。
retries: 0
#当有多个消息需要被发送到同一个分区时,生产者会把它们放在同一个批次里。该参数指定了一个批次可以使用的内存大小,按照字节数计算。
batch-size: 16384
# 设置生产者内存缓冲区的大小。
buffer-memory: 33554432
# 键的序列化方式
key-serializer: org.apache.kafka.common.serialization.StringSerializer
# 值的序列化方式
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# acks=0 : 生产者在成功写入消息之前不会等待任何来自服务器的响应。
# acks=1 : 只要集群的首领节点收到消息,生产者就会收到一个来自服务器成功响应。
# acks=all :只有当所有参与复制的节点全部收到消息时,生产者才会收到一个来自服务器的成功响应。
acks: 1
consumer:
# 自动提交的时间间隔 在spring boot 2.X 版本中这里采用的是值的类型为Duration 需要符合特定的格式,如1S,1M,2H,5D
auto-commit-interval: 1S
# 该属性指定了消费者在读取一个没有偏移量的分区或者偏移量无效的情况下该作何处理:
# latest(默认值)在偏移量无效的情况下,消费者将从最新的记录开始读取数据(在消费者启动之后生成的记录)
# earliest :在偏移量无效的情况下,消费者将从起始位置读取分区的记录
auto-offset-reset: earliest
# 是否自动提交偏移量,默认值是true,为了避免出现重复数据和数据丢失,可以把它设置为false,然后手动提交偏移量
enable-auto-commit: false
# 键的反序列化方式
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 值的反序列化方式
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
listener:
# 在侦听器容器中运行的线程数。
concurrency: 5
#listner负责ack,每调用一次,就立即commit
ack-mode: manual_immediate

生产者实例:

import com.xy.kafka.constant.Topic;
import com.xy.kafka.dao.XyKafkaInMsgMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.util.concurrent.ListenableFuture;
import java.util.Date;
import java.util.UUID;
@Slf4j
@Component
@EnableScheduling
public class MsgProducer {
@Autowired
private KafkaTemplate kafkaTemplate;
@Autowired
private XyKafkaInMsgMapper inMsgMapper;
/**
* 定时任务生成消息
*/
@Scheduled(cron = "0/10 * * * * ?")
public void send() {
//  要推送的消息内容
String message = UUID.randomUUID().toString()+"生产的消息";
/**
* 发送消息且接收返回值
*ListenableFuture 类是 spring 对 java 原生 Future 的扩展增强,是一个泛型接口,用于监听异步方法的回调而对于 kafka send 方法返回值而言,这里的泛型所代表的实际类型就是 SendResult<K, V>,而这里 K,V 的泛型实际上被用于 ProducerRecord<K, V> producerRecord,即生产者发送消息的 key,value 类型
*/
ListenableFuture listenableFuture = kafkaTemplate.send(Topic.SIMPLE,message);
listenableFuture.addCallback(
o -> log.info("消息发送成功,{}", o.toString()),
throwable -> log.info("消息发送失败,{}" + throwable.getMessage())
);
XyKafkaInMsg build = new XyKafkaInMsg();
build.setFwBh(System.currentTimeMillis());
build.setGmtCreate(new Date());
// 根据信息唯一批次号查询消息是否已存在
XyKafkaInMsg inMsg = inMsgMapper.selectByFuBh(build.getFwBh());
if (inMsg == null){
// 发送的消息入库
int saveMsgResult = inMsgMapper.insertSelective(build);
log.info("消息插入结果:{}",saveMsgResult == 1 ? "成功" : "失败");
}
}
}

消费者实例:

import com.xy.kafka.constant.Topic;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class KafkaSimpleConsumer {
// 简单消费者
@KafkaListener(groupId = "simpleGroup", topics = Topic.SIMPLE)
public void consumer1_1(ConsumerRecord<String, Object> record, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, Consumer consumer, Acknowledgment ack) {
log.info("单独消费者消费消息,topic= {} ,content = {}",topic,record.value());
log.info("consumer content = {}",consumer);
// 消息接收确认
ack.acknowledge();
/*
* 如果需要手工提交异步 consumer.commitSync();
* 手工同步提交 consumer.commitAsync()
*/
}
}

多消费者组的定义
https://blog.csdn.net/m0_3780…

javaspringspringboot
阅读 41更新于 今天 11:03
本作品系原创,采用《署名-非商业性使用-禁止演绎 4.0 国际》许可协议
avatar

一片秋叶一树春

贪君子之财,好美景之色,行正义之事,了前生之愿,爱此生之人!!!!!

37 声望
2 粉丝

0 条评论
得票时间

avatar

一片秋叶一树春

贪君子之财,好美景之色,行正义之事,了前生之愿,爱此生之人!!!!!

37 声望
2 粉丝

宣传栏

pom引入

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

yml配置说明

spring:
kafka:
# 以逗号分隔的地址列表,用于建立与Kafka集群的初始连接(kafka 默认的端口号为9092)
bootstrap-servers: 192.168.1.200:9092
producer:
# 发生错误后,消息重发的次数。
retries: 0
#当有多个消息需要被发送到同一个分区时,生产者会把它们放在同一个批次里。该参数指定了一个批次可以使用的内存大小,按照字节数计算。
batch-size: 16384
# 设置生产者内存缓冲区的大小。
buffer-memory: 33554432
# 键的序列化方式
key-serializer: org.apache.kafka.common.serialization.StringSerializer
# 值的序列化方式
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# acks=0 : 生产者在成功写入消息之前不会等待任何来自服务器的响应。
# acks=1 : 只要集群的首领节点收到消息,生产者就会收到一个来自服务器成功响应。
# acks=all :只有当所有参与复制的节点全部收到消息时,生产者才会收到一个来自服务器的成功响应。
acks: 1
consumer:
# 自动提交的时间间隔 在spring boot 2.X 版本中这里采用的是值的类型为Duration 需要符合特定的格式,如1S,1M,2H,5D
auto-commit-interval: 1S
# 该属性指定了消费者在读取一个没有偏移量的分区或者偏移量无效的情况下该作何处理:
# latest(默认值)在偏移量无效的情况下,消费者将从最新的记录开始读取数据(在消费者启动之后生成的记录)
# earliest :在偏移量无效的情况下,消费者将从起始位置读取分区的记录
auto-offset-reset: earliest
# 是否自动提交偏移量,默认值是true,为了避免出现重复数据和数据丢失,可以把它设置为false,然后手动提交偏移量
enable-auto-commit: false
# 键的反序列化方式
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 值的反序列化方式
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
listener:
# 在侦听器容器中运行的线程数。
concurrency: 5
#listner负责ack,每调用一次,就立即commit
ack-mode: manual_immediate

生产者实例:

import com.xy.kafka.constant.Topic;
import com.xy.kafka.dao.XyKafkaInMsgMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.util.concurrent.ListenableFuture;
import java.util.Date;
import java.util.UUID;
@Slf4j
@Component
@EnableScheduling
public class MsgProducer {
@Autowired
private KafkaTemplate kafkaTemplate;
@Autowired
private XyKafkaInMsgMapper inMsgMapper;
/**
* 定时任务生成消息
*/
@Scheduled(cron = "0/10 * * * * ?")
public void send() {
//  要推送的消息内容
String message = UUID.randomUUID().toString()+"生产的消息";
/**
* 发送消息且接收返回值
*ListenableFuture 类是 spring 对 java 原生 Future 的扩展增强,是一个泛型接口,用于监听异步方法的回调而对于 kafka send 方法返回值而言,这里的泛型所代表的实际类型就是 SendResult<K, V>,而这里 K,V 的泛型实际上被用于 ProducerRecord<K, V> producerRecord,即生产者发送消息的 key,value 类型
*/
ListenableFuture listenableFuture = kafkaTemplate.send(Topic.SIMPLE,message);
listenableFuture.addCallback(
o -> log.info("消息发送成功,{}", o.toString()),
throwable -> log.info("消息发送失败,{}" + throwable.getMessage())
);
XyKafkaInMsg build = new XyKafkaInMsg();
build.setFwBh(System.currentTimeMillis());
build.setGmtCreate(new Date());
// 根据信息唯一批次号查询消息是否已存在
XyKafkaInMsg inMsg = inMsgMapper.selectByFuBh(build.getFwBh());
if (inMsg == null){
// 发送的消息入库
int saveMsgResult = inMsgMapper.insertSelective(build);
log.info("消息插入结果:{}",saveMsgResult == 1 ? "成功" : "失败");
}
}
}

消费者实例:

import com.xy.kafka.constant.Topic;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class KafkaSimpleConsumer {
// 简单消费者
@KafkaListener(groupId = "simpleGroup", topics = Topic.SIMPLE)
public void consumer1_1(ConsumerRecord<String, Object> record, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, Consumer consumer, Acknowledgment ack) {
log.info("单独消费者消费消息,topic= {} ,content = {}",topic,record.value());
log.info("consumer content = {}",consumer);
// 消息接收确认
ack.acknowledge();
/*
* 如果需要手工提交异步 consumer.commitSync();
* 手工同步提交 consumer.commitAsync()
*/
}
}

多消费者组的定义
https://blog.csdn.net/m0_3780…

本文地址:H5W3 » 【Java】springboot2.x整合kafka

评论 0

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址