如何在Kafka Producer发送消息前验证JSON字段非空且age>0?
实现字段校验的几种方案
方案一:手动编写校验逻辑
最直接的方式是在发送消息前(比如sendMessage方法内部,或者Rest接口接收参数时)手动检查字段合法性:
1. 在sendMessage方法内添加校验
修改你的sendMessage方法,先校验Person对象的字段,不合法就抛出异常(或根据业务需求记录日志、返回错误响应):
public void sendMessage(Person p) { // 校验id非空且非空白 if (p.getId() == null || p.getId().trim().isEmpty()) { throw new IllegalArgumentException("id不能为空或空白"); } // 校验name非空且非空白 if (p.getName() == null || p.getName().trim().isEmpty()) { throw new IllegalArgumentException("name不能为空或空白"); } // 校验age大于0 if (p.getAge() == null || p.getAge() <= 0) { throw new IllegalArgumentException("age必须大于0"); } LOGGER.info(String.format("Message sent -> %s", p.toString())); Message<Person> message = MessageBuilder .withPayload(p) .setHeader(KafkaHeaders.TOPIC, "jsonTopic") .build(); kafkaTemplate.send(message); }
2. 在Rest接口层提前校验
如果是Spring Boot项目,建议在接收JSON的Controller层先做校验,避免无效对象流入Kafka发送环节:
@RestController @RequestMapping("/api") public class PersonController { private final KafkaProducerService producerService; public PersonController(KafkaProducerService producerService) { this.producerService = producerService; } @PostMapping("/person") public ResponseEntity<String> receivePerson(@RequestBody Person person) { // 提前校验参数 if (person.getId() == null || person.getId().trim().isEmpty()) { return ResponseEntity.badRequest().body("id不能为空或空白"); } if (person.getName() == null || person.getName().trim().isEmpty()) { return ResponseEntity.badRequest().body("name不能为空或空白"); } if (person.getAge() == null || person.getAge() <= 0) { return ResponseEntity.badRequest().body("age必须大于0"); } producerService.sendMessage(person); return ResponseEntity.ok("消息已接收并准备发送"); } }
方案二:使用Bean Validation(推荐)
如果是中大型Spring Boot项目,推荐用JSR-380标准的Bean Validation框架(比如Hibernate Validator),通过注解实现校验,逻辑更简洁易维护,还能复用校验规则。
1. 给Person实体类添加校验注解
import jakarta.validation.constraints.NotBlank; import jakarta.validation.constraints.Positive; public class Person { @NotBlank(message = "id不能为空或空白") private String id; @NotBlank(message = "name不能为空或空白") private String name; @Positive(message = "age必须大于0") private Integer age; // 构造器、getter、setter、toString方法省略 }
2. 在Rest接口层自动触发校验
在Controller方法的@RequestBody参数前加@Valid,Spring会自动校验参数,不合法时抛出MethodArgumentNotValidException,可以通过全局异常处理器统一返回错误响应:
@RestController @RequestMapping("/api") public class PersonController { private final KafkaProducerService producerService; public PersonController(KafkaProducerService producerService) { this.producerService = producerService; } @PostMapping("/person") public ResponseEntity<String> receivePerson(@Valid @RequestBody Person person) { producerService.sendMessage(person); return ResponseEntity.ok("消息已接收并准备发送"); } }
3. 全局异常处理(统一返回校验错误)
添加全局异常处理器,把校验失败的字段和错误信息整理成友好响应:
import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; import org.springframework.validation.FieldError; import org.springframework.web.bind.MethodArgumentNotValidException; import org.springframework.web.bind.annotation.ExceptionHandler; import org.springframework.web.bind.annotation.RestControllerAdvice; import java.util.HashMap; import java.util.Map; @RestControllerAdvice public class GlobalExceptionHandler { @ExceptionHandler(MethodArgumentNotValidException.class) public ResponseEntity<Map<String, String>> handleValidationExceptions(MethodArgumentNotValidException ex) { Map<String, String> errors = new HashMap<>(); ex.getBindingResult().getAllErrors().forEach((error) -> { String fieldName = ((FieldError) error).getField(); String errorMessage = error.getDefaultMessage(); errors.put(fieldName, errorMessage); }); return new ResponseEntity<>(errors, HttpStatus.BAD_REQUEST); } }
4. 在sendMessage方法内手动触发校验
如果需要在Kafka发送前再次校验(比如对象来自非Controller的其他途径),可以手动注入Validator触发校验:
import jakarta.validation.Validator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; import java.util.Set; @Service public class KafkaProducerService { private static final Logger LOGGER = LoggerFactory.getLogger(KafkaProducerService.class); private final KafkaTemplate<String, Person> kafkaTemplate; private final Validator validator; public KafkaProducerService(KafkaTemplate<String, Person> kafkaTemplate, Validator validator) { this.kafkaTemplate = kafkaTemplate; this.validator = validator; } public void sendMessage(Person p) { // 手动触发校验 Set<jakarta.validation.ConstraintViolation<Person>> violations = validator.validate(p); if (!violations.isEmpty()) { StringBuilder errorMsg = new StringBuilder(); for (jakarta.validation.ConstraintViolation<Person> violation : violations) { errorMsg.append(violation.getMessage()).append("; "); } throw new IllegalArgumentException("参数校验失败: " + errorMsg); } LOGGER.info(String.format("Message sent -> %s", p.toString())); Message<Person> message = MessageBuilder .withPayload(p) .setHeader(KafkaHeaders.TOPIC, "jsonTopic") .build(); kafkaTemplate.send(message); } }
选择建议
- 小型项目或简单场景,手动校验快速直接;
- 中大型项目推荐Bean Validation,校验逻辑与业务解耦,代码更简洁易维护。
内容的提问来源于stack exchange,提问作者Flavius Octavian
相关产品推荐
相关产品推荐

