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

迁移至Spring Data Elasticsearch:能否配置Ingest Pipeline用于文档索引?

在Spring Data Elasticsearch中配置Ingest Pipeline实现文档索引

完全可以在Spring Data Elasticsearch中为文档索引配置Ingest Pipeline,主要有以下几种实现方式:

1. 实体类全局绑定管道

通过@Document注解的pipeline属性,为指定实体的所有索引操作默认绑定Ingest Pipeline。这种方式适合所有该实体的索引请求都使用同一管道的场景:

@Document(indexName = "test-index", pipeline = "test-ingest-pipeline")
public class ContentTransformedModel {
    // 实体字段定义
}

之后直接使用Reactive Repository的save方法即可自动应用该管道:

public interface ContentTransformedRepository extends ReactiveElasticsearchRepository<ContentTransformedModel, String> {
}

// 业务代码
public Mono<ContentTransformedModel> addContent(RequestDoc requestDoc) {
    ContentTransformedModel model = convertToModel(requestDoc); // 自行实现请求到实体的转换逻辑
    return repository.save(model)
            .onErrorResume(IOException.class, e -> Mono.error(new ElasticsearchPublishException(ES_INDEX_INTERNAL_SERVER_ERROR, e)));
}

2. 单请求动态指定管道

如果需要针对特定请求灵活指定不同的Ingest Pipeline,可以借助ReactiveElasticsearchOperations(响应式场景)或ElasticsearchRestTemplate(同步场景),通过IndexQuery来配置:

响应式场景示例

@Autowired
private ReactiveElasticsearchOperations operations;

public Mono<IndexResponse> addContent(RequestDoc requestDoc) {
    ContentTransformedModel model = convertToModel(requestDoc);
    IndexQuery indexQuery = new IndexQueryBuilder()
            .withObject(model)
            .withIndexName("test-index")
            .withPipeline("test-ingest-pipeline")
            .build();

    return operations.index(indexQuery, IndexCoordinates.of("test-index"))
            .onErrorResume(IOException.class, e -> Mono.error(new ElasticsearchPublishException(ES_INDEX_INTERNAL_SERVER_ERROR, e)));
}

同步场景示例

@Autowired
private ElasticsearchRestTemplate restTemplate;

public IndexResponse addContent(RequestDoc requestDoc) {
    ContentTransformedModel model = convertToModel(requestDoc);
    IndexQuery indexQuery = new IndexQueryBuilder()
            .withObject(model)
            .withPipeline("test-ingest-pipeline")
            .build();
    try {
        return restTemplate.index(indexQuery);
    } catch (IOException e) {
        throw new ElasticsearchPublishException(ES_INDEX_INTERNAL_SERVER_ERROR, e);
    }
}

以上两种方式分别对应全局默认配置和单请求动态配置,能满足不同业务场景下的Ingest Pipeline使用需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 10:29:52