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}