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

使用Spring WebClient通过Elasticsearch _bulk API发送Flux遇400错误

问题分析与修复方案

你的代码返回HTTP 400错误,核心原因是请求体格式不符合Elasticsearch Bulk API要求,同时WebClient的请求体配置存在错误,以下是具体问题和修复步骤:

1. 核心错误:WebClient请求体类型不匹配

你构建的createCommandFlux是Flux<String>类型(每个元素是Bulk指令的NDJSON片段),但在WebClient的body()方法中错误指定了Foo.class:

.body(windowFlux, Foo.class)

这会导致WebClient尝试将字符串序列化为Foo对象,生成的请求体完全不符合NDJSON格式,直接触发Elasticsearch的400错误。

2. 修复步骤

步骤1:修正WebClient请求体配置

将body()方法的类型改为String.class,让WebClient直接发送字符串流:

.body(windowFlux, String.class)
// 或者更明确地使用BodyInserters:
.body(BodyInserters.fromPublisher(windowFlux, String.class))

步骤2:过滤无效的空字符串

当前代码在序列化失败时返回空字符串,会导致Bulk请求体中出现无效空行,引发Elasticsearch解析错误,需过滤掉空/无效条目:

Flux<String> createCommandFlux = Flux.interval(Duration.ofMillis(100))
        .map(i -> {
            try {
                Foo onePojo = new Foo(LocalDateTime.now().toString(), String.valueOf(i));
                String jsonStringOfOnePojo = new ObjectMapper().writeValueAsString(onePojo);
                // 每个Bulk条目由操作行+文档行组成,无需额外末尾换行(窗口拼接时会自动累加)
                return "{ \"create\" : {} }\n" + jsonStringOfOnePojo;
            } catch (Exception e) {
                e.printStackTrace();
                return null;
            }
        })
        .filter(Objects::nonNull); // 过滤无效条目

步骤3:优化窗口处理(可选)

单纯的window(100)在无限流场景下,最后一个未攒够100条的窗口会一直等待,无法触发请求。可以添加超时机制,确保数据及时发送:

.windowTimeout(100, Duration.ofSeconds(5)) // 每攒够100条或等待5秒,就触发一次批量请求

步骤4:添加响应错误处理

打印Elasticsearch返回的错误详情,方便调试问题:

.flatMap(clientResponse -> {
    if (clientResponse.statusCode().isError()) {
        return clientResponse.bodyToMono(String.class)
                .doOnNext(error -> System.err.println("Bulk请求失败:" + error))
                .flatMap(error -> Mono.error(new RuntimeException("Bulk API错误:" + error)));
    }
    return clientResponse.bodyToMono(String.class);
})

完整修复后的代码

import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.http.MediaType;
import org.springframework.web.reactive.function.BodyInserters;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.Disposable;
import reactor.netty.http.client.HttpClient;

import java.time.LocalDateTime;
import java.util.Objects;

public class BulkInsertOfFluxUsingWebClientBulkRestApi {

    public static void main(String[] args) throws InterruptedException {
        WebClient client = WebClient.create("http://127.0.0.1:9200/")
                .mutate()
                .clientConnector(new reactor.netty.http.client.ReactorClientHttpConnector(HttpClient.create().wiretap(true)))
                .build();

        ObjectMapper objectMapper = new ObjectMapper();
        Flux<String> createCommandFlux = Flux.interval(Duration.ofMillis(100))
                .map(i -> {
                    try {
                        Foo onePojo = new Foo(LocalDateTime.now().toString(), String.valueOf(i));
                        String jsonStringOfOnePojo = objectMapper.writeValueAsString(onePojo);
                        return "{ \"create\" : {} }\n" + jsonStringOfOnePojo;
                    } catch (Exception e) {
                        e.printStackTrace();
                        return null;
                    }
                })
                .filter(Objects::nonNull);

        Disposable disposable = createCommandFlux
                .windowTimeout(100, Duration.ofSeconds(5))
                .flatMap(windowFlux -> client
                        .post()
                        .uri("my_index/_bulk")
                        .contentType(MediaType.APPLICATION_NDJSON)
                        .body(BodyInserters.fromPublisher(windowFlux, String.class))
                        .exchangeToMono(clientResponse -> {
                            if (clientResponse.statusCode().isError()) {
                                return clientResponse.bodyToMono(String.class)
                                        .doOnNext(error -> System.err.println("Bulk请求失败响应:" + error))
                                        .flatMap(error -> Mono.error(new RuntimeException("Bulk API调用失败:" + error)));
                            }
                            return clientResponse.bodyToMono(String.class)
                                    .doOnNext(response -> System.out.println("Bulk请求成功响应:" + response));
                        }))
                .subscribe(
                        success -> {},
                        error -> System.err.println("订阅错误:" + error.getMessage())
                );

        Thread.sleep(1000000);
        disposable.dispose();
    }

    static class Foo {
        private String time;
        private String id;

        public Foo(String time, String id) {
            this.time = time;
            this.id = id;
        }

        // 必须添加getter,否则ObjectMapper无法序列化
        public String getTime() {
            return time;
        }

        public String getId() {
            return id;
        }
    }
}

额外注意事项

  • 确保Foo类有正确的getter方法,否则ObjectMapper无法将其序列化为JSON。
  • Elasticsearch Bulk API要求每个操作必须是两行:第一行是操作元数据(如{"create":{}}),第二行是文档内容,整个请求体需为换行分隔的JSON(NDJSON)。
  • 无限流场景下,windowTimeout比单纯的window更可靠,避免数据长时间累积无法发送。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 21:47:04