如何实现可接收批量记录的Spring Boot控制器?
创建Spring Boot批量Kafka记录接收端点
1. 定义数据模型
先创建匹配Confluent记录格式的POJO类,用于映射请求体数据:
// KafkaRecord.java public class KafkaRecord { private KafkaRecordValue value; // Getters & Setters public KafkaRecordValue getValue() { return value; } public void setValue(KafkaRecordValue value) { this.value = value; } } // KafkaRecordValue.java public class KafkaRecordValue { private String type; private String data; // Getters & Setters public String getType() { return type; } public void setType(String type) { this.type = type; } public String getData() { return data; } public void setData(String data) { this.data = data; } }
2. 配置JSON Lines消息转换器
Confluent示例用的是JSON Lines格式(每行一个独立JSON对象,非JSON数组),需要配置Spring支持这种媒体类型:
// WebConfig.java import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.context.annotation.Configuration; import org.springframework.http.converter.HttpMessageConverter; import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter; import org.springframework.web.servlet.config.annotation.WebMvcConfigurer; import java.util.List; @Configuration public class WebConfig implements WebMvcConfigurer { @Override public void configureMessageConverters(List<HttpMessageConverter<?>> converters) { MappingJackson2HttpMessageConverter converter = new MappingJackson2HttpMessageConverter(); // 添加JSON Lines媒体类型支持 converter.getSupportedMediaTypes().add(org.springframework.http.MediaType.parseMediaType("application/x-ndjson")); ObjectMapper mapper = new ObjectMapper(); // 允许将单行JSON解析为数组 mapper.enable(com.fasterxml.jackson.databind.DeserializationFeature.ACCEPT_SINGLE_VALUE_AS_ARRAY); converter.setObjectMapper(mapper); converters.add(converter); } }
3. 编写控制器端点
创建匹配指定路径的POST端点,处理批量记录并完成认证、数据解析逻辑:
// KafkaRecordController.java import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*; import java.util.List; @RestController @RequestMapping("/kafka/v3/clusters/{clusterId}/topics/{topicName}/records") public class KafkaRecordController { @PostMapping(consumes = "application/x-ndjson") public ResponseEntity<Void> receiveBatchRecords( @PathVariable String clusterId, @PathVariable String topicName, @RequestHeader("Authorization") String authorizationHeader, @RequestBody List<KafkaRecord> records) { // 验证Basic认证头 if (!validateBasicAuth(authorizationHeader)) { return ResponseEntity.status(HttpStatus.UNAUTHORIZED).build(); } // 遍历处理每条记录 for (KafkaRecord record : records) { String valueType = record.getValue().getType(); String base64Data = record.getValue().getData(); // 这里替换为你的业务逻辑:比如解码base64、发送到Kafka等 System.out.printf("Cluster: %s, Topic: %s, Type: %s, Data: %s%n", clusterId, topicName, valueType, base64Data); } return ResponseEntity.status(HttpStatus.CREATED).build(); } // 示例:Basic认证验证逻辑 private boolean validateBasicAuth(String authorizationHeader) { if (authorizationHeader == null || !authorizationHeader.startsWith("Basic ")) { return false; } // 实际场景中:解码base64凭证,与系统存储的密钥对比 String encodedCredentials = authorizationHeader.substring(6); // byte[] decoded = Base64.getDecoder().decode(encodedCredentials); // String[] credentials = new String(decoded).split(":"); return true; // 替换为真实验证逻辑 } }
4. 测试端点
用类似Confluent的cURL命令测试,注意指定Content-Type: application/x-ndjson:
curl -X POST -H "Content-Type: application/x-ndjson" \ -H "Authorization: Basic <BASE64-encoded-key-and-secret>" \ "http://localhost:8080/kafka/v3/clusters/my-cluster/topics/my-topic/records" \ -d '{"value": {"type": "BINARY", "data": "SGVsbG8gV29ybGQh"}} {"value": {"type": "BINARY", "data": "U3ByaW5nIEJvb3QgRXhwbG9yZXIh"}}'
补充说明
- 如果客户端只能发送JSON数组格式,可修改控制器
consumes为application/json,请求体改为[{"value": {...}}, {"value": {...}}]。 - 认证逻辑可根据需求替换为OAuth2、API Key等方式。
- 处理记录时,可集成Spring Kafka直接将数据转发到Kafka集群。
内容的提问来源于stack exchange,提问作者Manupriya Logus
相关产品推荐
相关产品推荐

