如何仅为Elasticsearch批量请求修改Socket超时时间?
我在使用spring-data-elasticsearch 7.10.0实现Elasticsearch批量索引时遇到了问题,处理约10GB的大数据流时偶尔触发Socket超时异常,错误信息如下:
Error while processing event: 60,000 milliseconds timeout on connection http-outgoing-3244 [ACTIVE]; nested exception is java.lang.RuntimeException: 60,000 milliseconds timeout on connection http-outgoing-3244 [ACTIVE]
我的批量索引实现代码:
@Override public final List<IndexedObjectInformation> bulkIndex(List<IndexQuery> queries, BulkOptions bulkOptions, IndexCoordinates index) { Assert.notNull(queries, "List of IndexQuery must not be null"); Assert.notNull(bulkOptions, "BulkOptions must not be null"); return bulkOperation(queries, bulkOptions, index); }
配置信息:
批量请求每批包含200条IndexQuery,无并发,全局ES配置如下:
spring: data: elasticsearch: max-connection-idle-time: 15000 # 15 seconds read-timeout: 7000 # 7 seconds socket-timeout: 60000 # 60 seconds connection-timeout: 4000 # 4 seconds
全局调高Socket超时能解决问题,但我希望仅针对批量索引请求调整超时,不影响读请求。我尝试修改传入的BulkOptions设置20分钟超时,但依然触发60秒超时。请问有没有办法仅为批量索引请求设置不同的Socket超时?
解决方案
你遇到的核心问题是:BulkOptions设置的是Elasticsearch服务端的批量请求处理超时,而报错是客户端的Socket超时(客户端等待服务端响应的最长时间),所以修改BulkOptions不会生效。下面是两种可行的针对性解决方案:
方案1:为批量请求单独创建RestHighLevelClient实例
创建一个独立的客户端实例,专门用于批量操作,配置更长的Socket/读超时,和默认客户端隔离:
@Configuration public class ElasticsearchConfig { // 默认客户端,用于读请求等常规操作 @Bean("defaultEsClient") public RestHighLevelClient defaultEsClient(ElasticsearchProperties properties) { ClientConfiguration clientConfiguration = ClientConfiguration.builder() .connectedTo(properties.getClusterNodes().toArray(new String[0])) .withConnectTimeout(properties.getConnectionTimeout()) .withSocketTimeout(properties.getSocketTimeout()) .withReadTimeout(properties.getReadTimeout()) .withMaxConnectionIdleTime(properties.getMaxConnectionIdleTime()) .build(); return RestClients.create(clientConfiguration).rest(); } // 批量操作专用客户端,设置20分钟超时 @Bean("bulkEsClient") public RestHighLevelClient bulkEsClient(ElasticsearchProperties properties) { ClientConfiguration clientConfiguration = ClientConfiguration.builder() .connectedTo(properties.getClusterNodes().toArray(new String[0])) .withConnectTimeout(properties.getConnectionTimeout()) .withSocketTimeout(1200000) // 20分钟 .withReadTimeout(1200000) .withMaxConnectionIdleTime(properties.getMaxConnectionIdleTime()) .build(); return RestClients.create(clientConfiguration).rest(); } }
然后在批量索引服务中注入这个专用客户端,自定义批量逻辑:
@Service public class BulkIndexService { private final RestHighLevelClient bulkEsClient; public BulkIndexService(@Qualifier("bulkEsClient") RestHighLevelClient bulkEsClient) { this.bulkEsClient = bulkEsClient; } public List<IndexedObjectInformation> customBulkIndex(List<IndexQuery> queries, BulkOptions bulkOptions, IndexCoordinates index) throws IOException { BulkRequest bulkRequest = new BulkRequest(); // 转换IndexQuery为BulkRequest操作 for (IndexQuery query : queries) { IndexRequest indexRequest = new IndexRequest(index.getIndexName()); if (query.getId() != null) { indexRequest.id(query.getId()); } indexRequest.source(query.getSource(), XContentType.JSON); bulkRequest.add(indexRequest); } // 设置服务端处理超时(对应BulkOptions的timeout) if (bulkOptions.getTimeout() != null) { bulkRequest.timeout(bulkOptions.getTimeout()); } // 执行批量请求 BulkResponse bulkResponse = bulkEsClient.bulk(bulkRequest, RequestOptions.DEFAULT); // 转换响应结果 List<IndexedObjectInformation> result = new ArrayList<>(); for (BulkItemResponse item : bulkResponse.getItems()) { result.add(new IndexedObjectInformation(item.getId(), item.getVersion(), item.isFailed())); } return result; } }
方案2:使用RequestOptions动态设置请求级超时
Spring Data Elasticsearch底层依赖RestHighLevelClient,可以通过自定义RequestOptions为单个批量请求设置超时:
先定义全局的批量专用RequestOptions:
private static final RequestOptions BULK_REQUEST_OPTIONS; static { RequestOptions.Builder builder = RequestOptions.DEFAULT.toBuilder(); // 设置客户端Socket和读超时为20分钟 builder.setSocketTimeout(1200000); builder.setReadTimeout(1200000); BULK_REQUEST_OPTIONS = builder.build(); }
然后扩展ElasticsearchRestTemplate,重写批量操作方法,使用自定义的RequestOptions:
public class CustomElasticsearchRestTemplate extends ElasticsearchRestTemplate { public CustomElasticsearchRestTemplate(RestHighLevelClient restHighLevelClient) { super(restHighLevelClient); } @Override protected <T> List<IndexedObjectInformation> bulkOperation(List<? extends Query> queries, BulkOptions bulkOptions, IndexCoordinates index, Class<T> clazz) { try { BulkRequest bulkRequest = createBulkRequest(queries, bulkOptions, index, clazz); // 使用自定义超时的RequestOptions执行请求 BulkResponse bulkResponse = getClient().bulk(bulkRequest, BULK_REQUEST_OPTIONS); return convertBulkResponse(bulkResponse); } catch (IOException e) { throw new ElasticsearchException("Bulk operation failed", e); } } }
将这个自定义模板注入Spring容器,替代默认的ElasticsearchRestTemplate,或者在批量操作时单独使用即可。
额外建议
除了调整超时,还可以优化批量请求的大小:比如适当减少每批文档数量(比如从200降到100),降低单批请求的数据量,从根源减少服务端处理时间,降低超时概率。
内容的提问来源于stack exchange,提问作者Aethernite

