Java实现CSV文件批量索引到Elasticsearch 7.12.0报错求助(替代Logstash方案)
解决Elasticsearch 7.12.0批量插入时的
no requests added报错 我来帮你拆解问题根源,一步步解决这个报错:
一、报错的直接原因:没有批量请求被添加到任务中
看你代码里的这段判断逻辑:
if (data.length == 3) { // 构建文档并添加批量请求的逻辑 }
但你的CSV样本每行是5个字段(id,profile_id,hier_name,attri_name,item),用逗号分割后data.length的值是5,这个条件永远不会成立,导致没有任何请求被加入到bulkRequest里,最终触发Validation Failed: 1: no requests added错误。
紧急修复这个逻辑:
把判断条件改成data.length ==5,同时修正字段映射(你代码里的字段名和CSV列不匹配),还要注意ES7.x已经废弃了索引类型(index type):
if (data.length == 5) { try { XContentBuilder xContentBuilder = jsonBuilder() .startObject() .field("id", data[0]) .field("profile_id", data[1]) .field("hier_name", data[2]) .field("attri_name", data[3]) .field("item", data[4]) .endObject(); // ES7.x移除了index type,prepareIndex只传索引名和文档ID即可 bulkRequest.add(client.prepareIndex(indexName) .id(data[0]) .setSource(xContentBuilder)); if ((count + 1) % 500 == 0) { count = 0; addDocumentToESCluser(bulkRequest, noOfBatch, count); noOfBatch++; // 注意:执行完批量请求后要重置bulkRequest,避免重复提交 bulkRequest = client.prepareBulk(); } } catch (Exception e) { e.printStackTrace(); } }
二、适配Elasticsearch 7.12.0的关键注意事项
彻底抛弃Index Type
ES7.x开始完全移除了索引类型的概念,你代码里的indexTypeName已经完全无效,继续使用会引发额外错误,必须从prepareIndex方法中移除这个参数。替换废弃的TransportClient
ES7.12.0中TransportClient已经标记为废弃,官方明确推荐使用RestHighLevelClient作为Java客户端,后续版本会彻底删除TransportClient。这里给你一个简化的RestHighLevelClient批量插入示例:
首先添加Maven依赖:
<dependency> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-high-level-client</artifactId> <version>7.12.0</version> </dependency>
然后修改客户端初始化和批量逻辑:
import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.HttpHost; import org.elasticsearch.action.bulk.BulkRequest; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.common.xcontent.XContentType; // 初始化Rest客户端 public RestHighLevelClient initRestClient() { return new RestHighLevelClient( RestClient.builder(new HttpHost("localhost", 9200, "http")) ); } // 批量导入CSV逻辑 public void CSVbulkImport(boolean isHeaderIncluded) throws IOException { BulkRequest bulkRequest = new BulkRequest(); File file = new File("/home/niteshb/Documents/workspace-spring-tool-suite-4-4.10.0.RELEASE/ElasticSearchService/src/main/resources/elasticdata.csv"); BufferedReader bufferedReader = new BufferedReader(new FileReader(file)); String line = null; int count = 0, noOfBatch = 1; if (isHeaderIncluded) { bufferedReader.readLine(); } while ((line = bufferedReader.readLine()) != null) { if (line.trim().isEmpty()) continue; String[] data = line.split(","); if (data.length == 5) { // 构建JSON文档(也可以用XContentBuilder) String jsonDoc = String.format( "{\"id\":\"%s\",\"profile_id\":\"%s\",\"hier_name\":\"%s\",\"attri_name\":\"%s\",\"item\":\"%s\"}", data[0], data[1], data[2], data[3], data[4] ); IndexRequest indexRequest = new IndexRequest(indexName) .id(data[0]) .source(jsonDoc, XContentType.JSON); bulkRequest.add(indexRequest); if ((count + 1) % 500 == 0) { // 执行批量请求 client.bulk(bulkRequest, RequestOptions.DEFAULT); bulkRequest = new BulkRequest(); // 重置批量任务 count = 0; System.out.println("Batch " + noOfBatch + " 导入完成"); noOfBatch++; } } else { System.out.println("无效数据行: " + line); } count++; } // 处理剩余的未提交文档 if (bulkRequest.numberOfActions() > 0) { client.bulk(bulkRequest, RequestOptions.DEFAULT); System.out.println("最后一批文档导入完成"); } bufferedReader.close(); }
三、额外优化建议
- 避免硬编码CSV字段长度,可以读取表头后动态映射字段,适配更多场景;
- 不要用
split(",")解析CSV,如果字段内容包含逗号会出错,建议用专门的CSV解析库(如OpenCSV); - 增加失败重试机制,对批量插入中失败的文档单独记录或重试。
内容的提问来源于stack exchange,提问作者user13593864
相关产品推荐
相关产品推荐

