SpringBoot 整合 Kafka

1、加入POM

1 <!-- https://mvnrepository.com/artifact/org.springframework.kafka/spring-kafka --> 2 <dependency> 3 <groupId>org.springframework.kafka</groupId> 4 <artifactId>spring-kafka</artifactId> 5 <version>2.6.4</version> 6 </dependency>

2、配置

1# kafka configuration 2spring.kafka.bootstrap-servers=110.121.233.203:9092 3# 大于0,生产者失败重试 4spring.kafka.producer.retries=1 5# 每次批发消息量 6spring.kafka.producer.batch-size=16384 7spring.kafka.producer.buffer-memory=33554432 8# 指定消息key和消息body的编解码 9spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer 10spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer 11spring.kafka.consumer.group-id=group_user

3、Application

1@EnableKafka 2@SpringBootApplication 3public class KafkaApplication { 4 5 public static void main(String[] args) { 6 Main.run(KafkaApplication.class, args); 7 } 8}

4、配置类

1@Configuration 2public class KafkaConfig { 3 4 @Bean 5 public NewTopic topic() { 6 return new NewTopic(Constants.TOPIC_USER, 1, (short) 1); 7 } 8}

5、其它测试代码

1public class Constants { 2 3 public static final String TOPIC_USER = "topic_user"; 4 5 public static final String GROUP_USER = "group_user"; 6} 7 8 9@Component 10public class UserProducer { 11 12 private final KafkaTemplate<String, String> kafkaTemplate; 13 14 public UserProducer(KafkaTemplate<String, String> kafkaTemplate) { 15 this.kafkaTemplate = kafkaTemplate; 16 } 17 18 public KafkaTemplate<String, String> getKafkaTemplate() { 19 return kafkaTemplate; 20 } 21}

消息生产者和消费者:

1@Api(value = "用户管理") 2@Log4j2 3@RestController 4@RequestMapping("/user") 5public class UserController { 6 7 private UserProducer userProducer; 8 private UserService userService; 9 10 public UserController(UserProducer userProducer, UserService userService) { 11 this.userProducer = userProducer; 12 this.userService = userService; 13 } 14 15 @ApiOperation(value = "保存或更新用户", notes = "userName是必输入项") 16 @ApiImplicitParams(value = { 17 @ApiImplicitParam(name = "id", value = "用户主键", dataTypeClass = Long.class, required = false), 18 @ApiImplicitParam(name = "userName", value = "用户姓名", dataTypeClass = String.class, required = true) 19 }) 20 @PostMapping("/saveOrUpdate") 21 public ResponseEntity<Result<User>> saveOrUpdate( 22 @RequestParam(required = false) Long id, 23 @RequestParam String userName 24 ) { 25 User user = new User(id, userName, new Date()); 26 userService.saveOrUpdate(user); 27 userProducer.getKafkaTemplate().send(Constants.TOPIC_USER, JSON.toJSONString(user).toString()); 28 return ResultUtils.ok(user); 29 } 30 31 @ApiOperation(value = "查询所有用户", notes = "") 32 @GetMapping("/queryAll") 33 public ResponseEntity<Result<List<User>>> queryAll() { 34 List<User> userList = userService.list(); 35 return ResultUtils.ok(userList, "查询成功"); 36 } 37} 38 39@ApiModel(value = "用户实体") 40@Data 41@NoArgsConstructor 42@AllArgsConstructor 43@TableName("t_base_user") 44public class User implements Serializable { 45 46 @ApiModelProperty(value = "主键") 47 @TableId(value = "id", type = IdType.NONE) 48 private Long id; 49 @ApiModelProperty(value = "姓名") 50 private String userName; 51 @ApiModelProperty(value = "创建时间") 52 private Date createTime; 53} 54 55@Log4j2 56@Component 57public class UserConsumer { 58 59 @KafkaListener(topics = {Constants.TOPIC_USER}, groupId = Constants.GROUP_USER) 60 public void onMessage(String message) { 61 log.error("[onMessage start]"); 62 log.error(message); 63 log.error("[onMessage end]"); 64 } 65}
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )