【问题标题】:Elastic Search Scroll API Asynchronous executionElastic Search Scroll API 异步执行
【发布时间】:2019-01-29 03:40:55
【问题描述】:

我正在运行具有 70Gb 索引大小/天的弹性搜索集群 5.6 版本。在一天结束时,我们被要求对过去 7 天的每个小时进行总结。我们正在使用 Java 版本的 High Level Rest 客户端,并且考虑到每个查询返回的文档数量对于滚动结果至关重要。

为了利用我们拥有的 CPU 并减少读取时间,我们正在考虑使用搜索滚动异步版本,但我们缺少一些示例,并且至少缺少其中的逻辑以继续前进。

我们已经检查了弹性相关文档,但它含糊不清:

https://www.elastic.co/guide/en/elasticsearch/client/java-rest/5.6/java-rest-high-search-scroll.html#java-rest-high-search-scroll-async

我们也在弹性讨论论坛中询问他们所说的但似乎没有人无法回答:

https://discuss.elastic.co/t/no-code-for-example-of-using-scrollasync-with-the-java-high-level-rest-client/165126

对此的任何帮助将不胜感激,并且可以肯定我不是唯一拥有此要求的人。

【问题讨论】:

    标签: api elasticsearch asynchronous search scroll


    【解决方案1】:

    这里是示例代码:

        public class App {
        public static void main(String[] args) throws IOException, InterruptedException {
            RestHighLevelClient client = new RestHighLevelClient(
                    RestClient.builder(HttpHost.create("http://localhost:9200")));
    
            client.indices().delete(new DeleteIndexRequest("test"), RequestOptions.DEFAULT);
            for (int i = 0; i < 100; i++) {
                client.index(new IndexRequest("test", "_doc").source("foo", "bar"), RequestOptions.DEFAULT);
            }
            client.indices().refresh(new RefreshRequest("test"), RequestOptions.DEFAULT);
    
            SearchRequest searchRequest = new SearchRequest("test").scroll(TimeValue.timeValueSeconds(30L));
            SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
            String scrollId = searchResponse.getScrollId();
    
            System.out.println("response = " + searchResponse);
    
            SearchScrollRequest scrollRequest = new SearchScrollRequest(scrollId)
                    .scroll(TimeValue.timeValueSeconds(30));
    
    
            //I was missing to wait for the results
            final CountDownLatch countDownLatch = new CountDownLatch(1);
    
            client.scrollAsync(scrollRequest, RequestOptions.DEFAULT, new ActionListener<SearchResponse>() {
                public void onResponse(SearchResponse searchResponse) {
                    System.out.println("response async = " + searchResponse);
                }
    
                public void onFailure(Exception e) {
    
                }
            });
    
            //Here we wait
            countDownLatch.await();
    
            //Clear the scroll if we finish before the time to keep it alive. Otherwise it will be clear when the time is reached.    
            ClearScrollRequest request = new ClearScrollRequest()
            request.addScrollId(scrollId);
    
            client.clearScrollAsync(request, new ActionListener<ClearScrollResponse>(){
               @Override
               public void onResponse(ClearScrollResponse clearScrollResponse) {
               }
    
               @Override
               public void onFailure(Exception e) {
               }
             });
    
            client.close();           
           }
        }
    

    感谢大卫·皮拉托 elastic discussion

    【讨论】:

      【解决方案2】:

      过去 7 天每小时的总结

      听起来您想对数据运行一些聚合,而不是获取原始文档。可能在第一级是日期直方图,以便以 1 小时的间隔进行汇总。在该日期直方图中,您需要一个内部 aggs 来运行您的聚合 - 指标/存储桶取决于所需的汇总。

      从 Elasticsearch v6.1 开始,您可以使用 Composite Aggregation 来使用分页获取所有结果存储桶。来自我链接的文档:

      复合聚合可用于有效地对多级聚合中的所有存储桶进行分页。这种聚合提供了一种流式传输特定聚合的所有存储桶的方法,类似于滚动对文档所做的操作。

      不幸的是,此选项在 v6.1 之前不存在,因此您要么需要升级 ES 才能使用它,要么找到另一种方式,例如拆分为多个查询,这样就可以满足 7 天的要求。

      【讨论】:

      • 谢谢@ziv,但我们希望 Elastic 搜索滚动 API 异步执行。我们需要对原始数据进行额外的操作。
      • 得到了一个高级客户端中scrollAsync的例子:github.com/elastic/elasticsearch/blob/… final CountDownLatch latch = new CountDownLatch(1); scrollListener = new LatchedActionListener(scrollListener,latch); client.scrollAsync(scrollRequest, RequestOptions.DEFAULT, scrollListener); // assertTrue(latch.await(30L, TimeUnit.SECONDS));
      • 是的,我们找到了同样的例子,仍然不是很清楚。考虑到这段代码@ziv,也许你可以帮助我们给出一个更好理解的例子
      • 就像client.scroll一样,client.scrollAsync只对ES运行一个请求。唯一的区别是请求是异步运行的,完成后你会收到一个回调调用。在响应中,您有 scrollId: scrollId = searchScrollResponse.getScrollId();然后,您可以使用 scrollId 来运行另一个异步滚动请求 - 这意味着迭代是通过使用更新后的 scrollId 再次调用 scrollAsync 的回调来完成的。我在 github 中看到的示例确实只显示了一个滚动调用,因此很难看到如何继续 - 通过回调
      猜你喜欢
      • 2017-05-17
      • 1970-01-01
      • 2016-06-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-08-03
      • 1970-01-01
      相关资源
      最近更新 更多