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
相关产品推荐
相关产品推荐

