使用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收到不支持的格式。
修改步骤如下:
- 替换
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") );
- 可以去掉手动设置的
Content-TypeHeader(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
相关产品推荐
相关产品推荐

