You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现可接收批量记录的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 19:57:48