Skip to main content
Version: 0.13.0

Concurrent segment search

Use concurrent segment search to search segments in parallel during the query phase. Concurrent segment search is enabled by default in Lucenia (search.concurrent_segment_search.mode is auto).

In auto mode, a shard request that includes an aggs section runs across segment slices when the shard is large enough. Aggregations visit every matching document on the shard. A request that only collects top hits (match, term, range, bool, knn) stays on a single thread.

auto aggregations-only behaviour, 0.13.0

Background​

In Lucenia, each search request follows the scatter-gather protocol. The coordinating node receives a search request, evaluates which shards are needed to serve this request, and sends a shard-level search request to each of those shards. Each shard that receives the request executes the request locally using Lucene and returns the results. The coordinating node merges the responses received from all shards and sends the search response back to the client. Optionally, the coordinating node can perform a fetch phase before returning the final results to the client if any document field or the entire document is requested by the client as part of the response.

Searching segments concurrently​

Without concurrent segment search, Lucene executes a request sequentially across all segments on each shard during the query phase. The query phase then collects the top hits for the search request. With concurrent segment search, each shard-level request will search the segments in parallel during the query phase. For each shard, the segments are divided into multiple slices. Each slice is the unit of work that can be executed in parallel on a separate thread, so the slice count determines the maximum degree of parallelism for a shard-level request. Once all the slices complete their work, Lucene performs a reduce operation on the slices, merging them and creating the final result for this shard-level request. Slices run on the index_searcher thread pool. Shard-level search requests still use the search thread pool.

Configuring the concurrent segment search mode​

Concurrent segment search is controlled by the search.concurrent_segment_search.mode dynamic cluster setting, which accepts the following values.

ValueDescription
auto(Default) Parallelize requests that include aggregations when the shard has enough segments and documents. Requests that only collect top hits stay on a single thread. See The auto decider.
allFan out across segment slices for every eligible request.
noneNever fan out; run every request single-threaded. This fully disables concurrent segment search and releases its resources.

The default is auto. To change the mode for all indexes in the cluster, set the dynamic cluster setting:

PUT _cluster/settings
{
"persistent":{
"search.concurrent_segment_search.mode": "all"
}
}

The auto decider​

In auto mode, a shard-level request runs in parallel across segment slices when all of the following are true:

  1. The search body includes an aggs (or aggregations) section. The whole query phase on that shard is sliced, including hit collection on the same request. A request whose body has no aggs section stays on one thread. That includes a long range and a knn top-k.
  2. The shard has at least two Lucene segments. Concurrent search splits work across those segments.
  3. The shard has at least search.concurrent_segment_search.auto.min_doc_count live documents (default 50000). That count is the live document count of the shard. A selective filter on a 200,000-document shard still fans out if the request has aggregations.

Confirm the settings:

GET _cluster/settings?include_defaults=true&filter_path=**.search.concurrent_segment_search*

To raise the live-document floor:

PUT _cluster/settings
{
"persistent":{
"search.concurrent_segment_search.auto.min_doc_count": 100000
}
}

To parallelize requests that collect only top hits (keyword search, knn, range without aggs), once they meet the segment and document floors:

PUT _cluster/settings
{
"persistent":{
"search.concurrent_segment_search.auto.fan_out_non_aggregating": true
}
}

The default is false. Requests that include aggregations fan out under both values of fan_out_non_aggregating.

Keep mode at auto and fan_out_non_aggregating at false for a cluster that runs aggregations and top-k search together. Set mode: all to slice every eligible request, including top-k.

Back-compatibility boolean setting​

The boolean search.concurrent_segment_search.enabled cluster setting (default true) and its index-level counterpart index.search.concurrent_segment_search.enabled remain for older configurations. Prefer search.concurrent_segment_search.mode for new clusters.

Lucenia evaluates the boolean after the mode check:

  • With mode set to none, concurrent segment search is off for every index. The cluster and index booleans are not consulted.
  • With mode set to auto or all, Lucenia reads index.search.concurrent_segment_search.enabled when that setting is present on the index; otherwise it uses search.concurrent_segment_search.enabled. A value of false keeps that index on a single thread for both auto and all. A value of true allows the mode decider to run for that index.

Leave the index-level boolean unset to inherit the cluster-level boolean. Retrieve the current index value with the Index Settings API and omit ?include_defaults; the setting appears only when it is explicitly set.

To toggle the boolean setting for a particular index, specify the index name in the endpoint:

PUT <index-name>/_settings
{
"index.search.concurrent_segment_search.enabled": true
}

Settings​

SettingScopeTypeDefaultDescription
search.concurrent_segment_search.modeCluster (dynamic)string (none, auto, all)autoWhen a shard-level request fans out across segment slices.
search.concurrent_segment_search.auto.min_doc_countCluster (dynamic)integer50000In auto mode, minimum live documents on the shard before fan-out.
search.concurrent_segment_search.auto.fan_out_non_aggregatingCluster (dynamic)BooleanfalseIn auto mode, whether requests without aggregations may fan out.
search.concurrent.max_slice_countCluster (dynamic)integer0Maximum slices per shard-level request. 0 uses Lucene slicing (up to 250K documents or 5 segments per slice). A positive integer uses round-robin max-slice-count slicing.
search.concurrent_segment_search.enabledCluster (dynamic)BooleantrueLegacy cluster gate used when mode is auto or all.
index.search.concurrent_segment_search.enabledIndex (dynamic)Booleanunset (inherits cluster)When present, replaces the cluster boolean for that index. Ignored when mode is none.

These cluster settings are also listed under Search settings. The index setting is listed under Dynamic index-level index settings.

The index_searcher thread pool​

CAT reports the pool name as index_searcher. The pool type is RESIZABLE. Each segment slice is one task on this pool. The size and queue_size fields are not in the default CAT columns; request them:

GET _cat/thread_pool/index_searcher?v&h=node_name,name,type,size,queue_size,active,queue,rejected
SettingDefaultDescription
thread_pool.index_searcher.sizeallocatedProcessorsMaximum concurrent slice tasks on the node. Node-scoped; set in lucenia.yml and restart the node to change it.
thread_pool.index_searcher.queue_size1000Queue depth for slice tasks waiting for a free index_searcher thread. Dynamic node setting.

Slice count is the maximum parallelism for one shard-level request. Pool size is the maximum number of slice tasks the node runs at once across all concurrent searches. When slice count times in-flight shard requests exceeds pool size, additional slices wait in the queue up to queue_size.

Slicing mechanisms​

You can choose one of two available mechanisms for assigning segments to slices: the default Lucene mechanism or the max slice count mechanism.

The Lucene mechanism​

By default, Lucene assigns a maximum of 250K documents or 5 segments (whichever is met first) to each slice in a shard. For example, consider a shard with 11 segments. The first 5 segments have 250K documents each, and the next 6 segments have 20K documents each. The first 5 segments will be assigned to 1 slice each because they each contain the maximum number of documents allowed for a slice. Then the next 5 segments will all be assigned to another single slice because of the maximum allowed segment count for a slice. The 11th slice will be assigned to a separate slice.

The max slice count mechanism​

The max slice count mechanism is an alternative slicing mechanism that uses a dynamically configurable maximum number of slices and divides segments among the slices in a round-robin fashion. This is useful when there are already too many top-level shard requests and you want to limit the number of slices per request in order to reduce competition between the slices.

Setting the slicing mechanism​

By default, concurrent segment search uses the Lucene mechanism to calculate the number of slices for each shard-level request. To use the max slice count mechanism instead, configure the search.concurrent.max_slice_count cluster setting:

PUT _cluster/settings
{
"persistent":{
"search.concurrent.max_slice_count": 2
}
}

The search.concurrent.max_slice_count setting can take the following valid values:

  • 0: Use the default Lucene mechanism.
  • Positive integer: Use the max target slice count mechanism. Usually, a value between 2 and 8 should be sufficient.

Limitations​

The following aggregations do not support the concurrent search model. If a search request contains one of these aggregations, the request will be executed using the non-concurrent path even if concurrent segment search is enabled at the cluster level or index level.

  • Parent aggregations on join fields.
  • sampler and diversified_sampler aggregations.
  • Composite aggregations that use scripts. Composite aggregations without scripts do support concurrent segment search.

Other considerations​

The following sections provide additional considerations for concurrent segment search.

The terminate_after search parameter​

The terminate_after search parameter is used to terminate a search request once a specified number of documents has been collected. If you include the terminate_after parameter in a request, concurrent segment search is disabled and the request is run in a non-concurrent manner.

Typically, queries are used with smaller terminate_after values and thus complete quickly because the search is performed on a reduced dataset. Therefore, concurrent search may not further improve performance in this case. Moreover, when terminate_after is used with other search request parameters, such as track_total_hits or size, it adds complexity and changes the expected query behavior. Falling back to a non-concurrent path for search requests that include terminate_after ensures consistent results between concurrent and non-concurrent requests.

Sorting​

Requests that sort on a time-series field always run in a non-concurrent manner, regardless of the configured mode, so that the time-series sort optimization is preserved.

For other sorts, the behavior depends on the data layout of the segments. The sort optimization feature can prune entire segments based on the min and max values as well as previously collected values. If the top values are present in the first few segments and all other segments are pruned, query latency may increase when sorting with concurrent segment search. Conversely, if the last few segments contain the top values, then latency may improve with concurrent segment search.

Terms aggregations​

Non-concurrent search calculates the document count error and returns it in the doc_count_error_upper_bound response parameter. During concurrent segment search, the shard_size parameter is applied at the segment slice level. Because of this, concurrent search may introduce an additional document count error.

Developer information: AggregatorFactory changes​

Because of implementation details, not all aggregator types can support concurrent segment search. To accommodate this, we have introduced a supportsConcurrentSegmentSearch() method in the AggregatorFactory class to indicate whether a given aggregation type supports concurrent segment search. By default, this method returns false. Any aggregator that needs to support concurrent segment search must override this method in its own factory implementation.

To ensure that a custom plugin-based Aggregator implementation works with the concurrent search path, plugin developers can verify their implementation with concurrent search enabled and then update the plugin to override the supportsConcurrentSegmentSearch() method to return true.