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

Spring Data Elasticsearch百万级数据批量索引/删除性能优化咨询

你的思路是对的,单条操作频繁发起ES请求IO开销过高,全量缓存所有请求再一次性批量执行确实会出现内存溢出,采用分批批量处理即可兼顾性能和内存占用,以下是具体落地方案:

索引场景优化

核心逻辑是遍历数据时攒批,每达到指定批次阈值就执行一次bulk写入,写完清空批次列表释放内存,内存仅会占用单批次的大小,不会出现OOM问题。
改造后代码示例:

// 批次大小可根据实际文档大小调整,推荐单批次1000~3000条,总大小控制在5~15MB(ES官方推荐最优区间)
private static final int BULK_BATCH_SIZE = 1000;

// 原遍历逻辑改造
List<IndexQuery> batchQueries = new ArrayList<>(BULK_BATCH_SIZE);
IndexCoordinates indexCoordinates = IndexCoordinates.of("portal_idx");
for (Map.Entry<Integer, ObjectDetails> key : objectDetailsHashMap.entrySet()) {
    IndexQuery indexQuery = buildIndexQuery(key, oPath);
    batchQueries.add(indexQuery);
    // 达到批次阈值执行批量写入
    if (batchQueries.size() >= BULK_BATCH_SIZE) {
        elasticsearchOperations.bulkIndex(batchQueries, indexCoordinates);
        // 清空批次列表释放内存
        batchQueries.clear();
    }
    // 其他数据库表数据插入逻辑...
}
// 遍历结束后处理剩余不满一批的数据
if (!batchQueries.isEmpty()) {
    elasticsearchOperations.bulkIndex(batchQueries, indexCoordinates);
}

// 原indexDocument方法改造为仅构造IndexQuery,不执行单条写入
private IndexQuery buildIndexQuery(Map.Entry<Integer, ObjectDetails> key, String oPath) {
    String docId = "" + key.getValue().getCatalogId() + key.getValue().getObjectId();

    byte[] nameBytes = key.getValue().getName();
    byte[] physicalNameBytes = key.getValue().getPhysicalName();
    byte[] definitionBytes =  key.getValue().getDefinition();
    byte[] commentBytes = key.getValue().getComment();

    return new IndexQueryBuilder()
            .withId(docId)
            .withObject(new MetadataSearch(
                    key.getValue().getObjectId(),
                    key.getValue().getCatalogId(),
                    key.getValue().getParentId(),
                    key.getValue().getTypeCode(),
                    key.getValue().getStartVersion(),
                    key.getValue().getEndVersion(),
                    nameBytes != null ? new String(nameBytes, StandardCharsets.UTF_8) : "-",
                    physicalNameBytes != null ? new String(physicalNameBytes, StandardCharsets.UTF_8) : "-",
                    definitionBytes != null ? new String(definitionBytes, StandardCharsets.UTF_8) : "-",
                    commentBytes != null ? new String(commentBytes, StandardCharsets.UTF_8) : "-",
                    oPath
            ))
            .build();
}

该方案相比原单条写入性能可以提升数倍到数十倍。

删除场景优化

有两种方案可选,可根据业务场景适配:

方案1:分批批量删除(适合已有待删除ID列表的场景)

和索引逻辑一致,攒够一批ID就执行bulkDelete,改造后代码:

private static final int DELETE_BATCH_SIZE = 1000;

private void deleteElasticDocuments(String catalogId) {
    String queryText = martServerContext.getQueryCacheInstance().getQuery(QUERY_PORTAL_GET_OBJECTS_IN_PORTAL_BY_MODEL);
    MapSqlParameterSource mapSqlParameterSource = new MapSqlParameterSource();
    mapSqlParameterSource.addValue("cId", Integer.parseInt(catalogId));
    List<String> deleteIds = new ArrayList<>(DELETE_BATCH_SIZE);
    IndexCoordinates indexCoordinates = IndexCoordinates.of("portal_idx");
    namedParameterJdbcTemplate.query(queryText, mapSqlParameterSource, (resultSet -> {
        int objectId = resultSet.getInt(O_ID);
        String docId = catalogId + objectId;
        deleteIds.add(docId);
        if (deleteIds.size() >= DELETE_BATCH_SIZE) {
            elasticsearchOperations.bulkDelete(deleteIds, indexCoordinates);
            deleteIds.clear();
        }
    }));
    // 处理剩余待删除ID
    if (!deleteIds.isEmpty()) {
        elasticsearchOperations.bulkDelete(deleteIds, indexCoordinates);
    }
}

方案2:按查询直接删除(适合按固定维度删除的场景,性能更高)

你当前的场景是删除指定catalogId下的所有文档,完全不需要先查数据库捞ID再删,直接构造ES查询条件删除所有匹配的文档,一步到位,性能远高于先查再删的方案,代码示例:

import org.springframework.data.elasticsearch.core.query.NativeSearchQuery;
import org.springframework.data.elasticsearch.core.query.NativeSearchQueryBuilder;
import static org.elasticsearch.index.query.QueryBuilders.termQuery;

private void deleteElasticDocuments(String catalogId) {
    // 构造查询条件:匹配catalogId等于传入值的所有文档
    NativeSearchQuery deleteQuery = new NativeSearchQueryBuilder()
            .withQuery(termQuery("catalogId", catalogId))
            .build();
    // 直接执行按查询删除
    elasticsearchOperations.delete(deleteQuery, MetadataSearch.class, IndexCoordinates.of("portal_idx"));
}

如果待删除文档量级超过千万级,可以给查询加上滚动分页参数,避免单次删除请求超时。

附加调优建议
  • 批次大小不要固定写死,可根据单文档平均大小测试调整,单批量请求总大小控制在5MB~15MB区间性能最优。
  • 如果是全量离线同步索引的场景,可以临时把索引的refresh_interval设置为-1,同时关闭索引副本,写入完成后再恢复原有配置,写入速度可以提升3~5倍。实时同步场景不要修改该配置。
  • 批量操作建议增加返回结果校验,捕获批量操作中的失败文档,记录日志后重试,避免数据不一致。
  • 业务允许的情况下,可以将批量操作设置为异步执行,避免阻塞ETL主流程。

内容的提问来源于stack exchange,提问作者this.srivastava

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 10:00:06