Solr/Elasticsearch capacity planning

When it comes to estimating resources needed for your index, there is no precise formula that will tell you how many shards you need to split data to or how many nodes you need to host it. It is a complex issue with a lot of variables and often you have less “equations” and there might be more than one right solution and even more often each solution comes with some compromises. However, there are some steps you can take to add some more inputs and help you choose setup that is more likely to perform well. You have to start from known requirements and limitations.

Index requirements

One way or another, you need to answer the following questions as precisely as possible:

  • What is the expected number of document in your index?
  • What is the maximum indexing throughput that you need to handle? Make sure that you use maximum and not average as you have to account for the worst case scenario.
  • What is the maximum time allowed between document being indexed and visible in search results, a.k.a near real time (NRT) requirement?
  • What sort of queries are expected and at what rate? Similar to indexing rate, maximum value should be used.
  • What is the average and maximum query latency allowed for different query groups?

If you are replacing some existing system, and if you are using some monitoring tool you can easily answer the previous questions. One such tool is Sematext Cloud that provides charts for all required metrics.

Indexing rate chart from Sematext Cloud

If you are starting with a new project, this step can be challenging. In some cases some metrics can be calculated - e.g. if you are indexing some periodic events from N devices. In other cases, you will have to do your best guess and later verify and adjust.

Known limitations

There are a few hard limitations that need to be considered:

  • Shard cannot contain more than 2 billion documents since Lucene is using integer for internal IDs. It is unlikely that you will reach this limitation without hitting some other limitation first.
  • Index size on disk should not take the entire disk because of ongoing merges. In theory, it should take only a half of it, but in practice it can be more if you are not planning on doing force merges to a single segment. In case you plan on having multiple tiers and have some indices that will become static you can use more disk space since merges are happening only while indexing.

Estimating unknown

Let’s assume you came up with following numbers:

  • You will have ~20M documents in index.
  • Max indexing rate at any moment will not be more than 1K docs/s.
  • Max number of concurrent users is 10 (it is also common to express this requirement as max query rate).

Before running any tests, referent node spec and heap size should be determined. Let’s assume we start with some smaller instance with 4 CPUs, 16GB RAM and 500GB of disk space and decide to start JVM with 8GB heap.

Tests should be performed, per index, in order to:

  • Determine max shard size allowed in order to have acceptable query latency - assume tests showed it is 500K docs and that such shard takes 5GB on disk.
  • Determine max indexing throughput per node/shard - assume tests showed that a single node can index up to 500 docs/s.
  • Determine max data per node to have optimal resource utilisation
  • Determine max possible query rate per node

Max shard size determines min number of shards:

shards_count ≥ round_up(max_docs / max_shard_size)
In our example, expected data volume is ~20M and max shard size is 500K documents, so minimum number of shards is 40.

It needs to be combined with max indexing throughput per node/shard in order to determine min number of nodes:

nodes_count ≥ round_up(max_indexing_rate / node_indexing_throughput)
For our numbers, max indexing throughput per node is 500 docs/s and since we expect no more than 1K docs/s indexing rate, 2 nodes should be enough to handle indexing load.

We continue with other limitations such as storage capacity to limit max shards per node:

shards_per_node ≤ round_down(disk_capacity / shard_size)
In our case, max shards per node is 10, but note that it needs to be smaller to handle merges.

Optimal number of shards per node depends also on query load. Concurrent queries can lead to CPU and memory starvation in case too many shards are on the same node. On the other hand, a single shard per node will result in max query rate, but resources utilisation might be low.

Once optimal number of shards per node is determined, min number of nodes can be calculated:

nodes_count ≥ round_up(shards_count / max_shards_per_node)
Let’s assume tests showed that 5 shards can be safely run per node. Min number of shards is already estimated to 40, so we can conclude that at least 8 nodes are needed in order to serve a single copy of shards.

Tests needs to be done in order to determine max query rate per node. Let’s assume that tests showed that 4 concurrent users can be served with a single set of shards (this is usually slightly higher number than number of cores that are available for processing queries). Since we need 10 concurrent users, 2 replicas (3 copies) are needed in order to support this requirement, so number of nodes should be multiplied by 3. Even if tests show that a single copy can satisfy search requirements, we need to have at least one replica in order to satisfy FT and HA requirements.

This is an iterative process. Node configuration and ES heap size could be revisited in order to see if some resource is under-utilised or if small changes could reduce costs. Whenever some estimate or node spec is changed, process needs to be redone, until optimal solution is found.

Note that the initial number of nodes does not have to be the same as the estimated number of nodes if initial data size is not the same as one used in tests. What usually counts is data volume that node serves. While shards are still small, more shards can be placed per node. As data grows and resource utilisation gets closer to max allowed, new node(s) can be added to reduce load per node.

What needs to be considered is whether new node(s) can lead to unbalanced cluster where one node serves significantly more data. That should not be the case with large enough number of shards and shards that are not too big. For smaller number of shards, preferred number of shards is a multiple of 12 because it will let us add a single node and still keep cluster balanced.

Similar logic can be applied to needed query rate - if max expected query rate is still not reached, there is no need to start with all the replicas.

Running tests

Now that we know how to combine test results, all we have to do is run the tests and obtain those numbers. Even though tests are described separately they are usually run in an iterative approach, combined and adjusted based on found results.

Max shard size

Shard size determines query latency - the bigger the shard, the slower the queries. In order to determine max shard size we run various types of queries while increasing shard size. Query time is measured observed from client’s perspective (include query/fetch/network latency) because that is how requirements are usually expressed. Since we are looking for max possible shard size, each query type can be run separately, without any concurrency, without any, or with min indexing. More complex queries can be run first since it is more likely to hit shard size limit with such queries. Once the size is determined, it can be used to determine min number of shards and the expected update rate per shard. Queries should be rerun with expected shard indexing rate to see how indexing affects query time. Shard size should be adjusted in order to meet query latency requirement. At this moment it is worth revisiting NRT requirement since it can significantly affect both query latency and indexing rate.

Max indexing throughput per node

Max indexing rate should not be determined assuming all node’s resources will be used for indexing, but in context of our query requirements. In case of index that neither indexing nor query dominate, we can assign 2 CPUs to indexing. Indexing requires CPU for analysis and segment merges (segments are compressed so CPU is needed for decompression). We can limit resources used for indexing by setting bulk thread pool in ES but it is better if we do that by limiting the number of indexing threads on client. Number of merging threads should also be tuned to control the load it can produce. Tests should run long enough to see throughput during major GC and large segment merges. Different bulk sizes and number of indexing threads should be used to see max throughput.

Refresh interval can increase indexing throughput significantly. Some tests suggest that changing refresh/commit interval from 1s to 5s increase throughput by 25% while increasing to 30s increases throughput by 75% so it should be set to max possible value.

Max shards per node

Single shard per node will make sure all resources are available for that shard, but that can lead to resources being under-utilised. Reducing node spec is an option, but that can lead to large number of nodes. It is more common to have more than one shard per node. Combined query/indexing test should be run to determine max query load when a single shard is hosted. While indexing load per shard is fixed, the number of concurrent users should be increased until latency is acceptable. After setting a baseline, new shard can be added, along with its indexing rate. Query rate should be split between shards and number of concurrent users adjusted. With each new shard, indexing is favorised and query penalised. There cannot be more shards than max allowed indexing rate per node divided by needed indexing rate per shard. In some cases instead of reducing the number of concurrent users it makes sense to reduce shard size. That can trigger changes in indexing rate per shard and tests need to be rerun.

Max number of concurrent users

This is partially done with the previous test. What needs to be additionally addressed is distributed nature of ES queries. There has to be room for merging part of distributed queries.

In case document routing is used, in needs to be included when calculating overall max number of concurrent users.

Conclusion

Estimating hardware requirements and number of shards and replicas is just an estimate. It needs to be verified using production load. Note that these are just guidelines that could work for some simpler usecases and it needs to be adjusted to particular cases - e.g. in case of static index, focus should be on queries, while some logging case should put stress on indexing. Hopefully this will give you ideas and starting point and you will be kind enough to share your findings in comments.

Post a Comment