I am following the AWS documentation for "Choosing the number of shards" for an Elasticsearch Index.
My Read TPS for the ES Index will be very high (around 1300 TPS, and can increase to 6500 TPS), but the amount of data which will be present will be very less (lesser than a GB).
- To match with the high reads, I am planning to implement horizontal scaling (increase the number of data nodes)
- Due to the very less data, as per the above documentation, the number of shards should be 1 (optimal desired shard size ~ 10GB-50GB, and my data being less than 1 GB)
Questions:
- As far as I can understand, one shard is not distributed over data nodes. (One shard can reside only on one data node). Is this understanding correct?
- From here,
In Elasticsearch, each query is executed in a single thread per shard. Multiple shards can however be processed in parallel, as can multiple queries and aggregations against the same shard.. If the above understanding is correct, all the requests will be single threaded on a single data node, if I only have one shard. The horizontal scaling thus cannot be implemented.
What should be the optimal number of primary shards/replicas for an index given the high TPS and low data?
Should I- still have a single shard, but multiple replicas (proportional to the number of hosts), or
- multiple primary shards itself (whose size would be in MBs), and a single replica (to save on the memory). (I don't see nodes going down in my cluster that badly that I NEED more than one replica!)