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

使用Kafka推送数据至OpenSearch时持续收到406状态码求助

问题描述

我正在编写基于Kafka的RESTful服务,要把Kafka消费者接收的消息存入OpenSearch索引,但每次发起请求时都收到以下错误:

Result: {"error":"Content-Type header [application/x-www-form-urlencoded] is not supported","status":406}

我的Kafka消费者代码如下:

public class Main {
    static Client client = new Client();

    public static void main(String[] args) throws IOException, NoSuchAlgorithmException, KeyManagementException {

        String servers = "localhost:9092";
        String groupId = "id";
        String topic = "mytopic";

        Properties properties = new Properties();
        properties.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
        properties.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<>(properties);
        kafkaConsumer.subscribe(Arrays.asList(topic));

        Logger logger = LoggerFactory.getLogger(Main.class.getName());

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(3000));
            for (ConsumerRecord<String, String> record : records) {
                logger.info("Key:" + record.key() + ", Value:" + record.value());
                logger.info("Partition:" + record.partition() + ", Offset:" + record.offset());

                String result = client.run("https://localhost:9200/index3/_create/" + record.offset(), "admin", "admin", record);
                System.out.println("Result: " + result);
            }
        }
    }
}

消息内容是"Request":"valid",我试过写成键值对JSON格式但没成功。后来尝试把Content-Type设置为application/json,但还是抛出相同错误!

我的客户端代码如下:

public String run(String url, String username, String password, ConsumerRecord<String, String> record) throws IOException, NoSuchAlgorithmException, KeyManagementException {

    //TrustManager per accettare tutti i certificati
    TrustManager[] trustAllCerts = new TrustManager[]{new Client()};

    SSLContext sslContext = SSLContext.getInstance("TLS");
    sslContext.init(null, trustAllCerts, new SecureRandom());

    //Client
    OkHttpClient client = new OkHttpClient.Builder()
            .connectionSpecs(Arrays.asList(ConnectionSpec.MODERN_TLS, ConnectionSpec.COMPATIBLE_TLS))
            .sslSocketFactory(sslContext.getSocketFactory(), (X509TrustManager) trustAllCerts[0])
            .hostnameVerifier((hostname, session) -> true)
            .build();

    //Authentication
    String credentials = username + ":" + password;
    String encode = Base64.getEncoder().encodeToString(credentials.getBytes());
    String authHeader = "Basic " + encode;


    RequestBody requestBody = new FormBody.Builder()
            .add("Message:", "n°" + record.offset())
            .build();

    Request request = new Request.Builder()
            .url(url)
            .header("Authorization", authHeader)
            .header("Content-Type", "application/json")
            .post(requestBody) //Body della post
            .build();

    try (Response response = client.newCall(request).execute()) {
        return response.body().string();
    }
}
解决方案

问题核心是RequestBody的类型冲突:你用FormBody.Builder构造请求体时,OkHttp会自动将请求的Content-Type设为application/x-www-form-urlencoded,就算手动设置Content-Type: application/json也会被覆盖,导致OpenSearch收到不支持的格式。

修改步骤如下:

  1. 替换FormBody为JSON格式的RequestBody,明确指定媒体类型为JSON:
// 构造包含Kafka消息和偏移量的JSON请求体
String jsonBody = String.format("{\"Message\":\"n°%d\",\"Request\":\"%s\"}", record.offset(), record.value());

RequestBody requestBody = RequestBody.create(
    jsonBody,
    MediaType.parse("application/json; charset=utf-8")
);
  1. 可以去掉手动设置的Content-Type Header(RequestBody已指定MediaType,OkHttp会自动添加正确的Header),保留也不会冲突。

修改后的客户端run方法关键部分:

//Authentication
String credentials = username + ":" + password;
String encode = Base64.getEncoder().encodeToString(credentials.getBytes());
String authHeader = "Basic " + encode;

// 构造JSON请求体
String jsonBody = String.format("{\"Message\":\"n°%d\",\"Request\":\"%s\"}", record.offset(), record.value());
RequestBody requestBody = RequestBody.create(
    jsonBody,
    MediaType.parse("application/json; charset=utf-8")
);

Request request = new Request.Builder()
        .url(url)
        .header("Authorization", authHeader)
        .post(requestBody)
        .build();

额外说明:如果Kafka消息本身就是合法JSON(比如{"Request":"valid"}),可以直接用消息内容作为请求体:

RequestBody requestBody = RequestBody.create(
    record.value(),
    MediaType.parse("application/json; charset=utf-8")
);

这样修改后,请求的Content-Type会正确设置为application/json,符合OpenSearch的要求,即可解决406错误。


内容的提问来源于stack exchange,提问作者user20174493

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:33:08