8 ms·
Mostly good stuff but a few comments: - article doesn’t clarify if it’s on hardware or VMs - 140 shards per node is certainly on the low side, one can easily
by DmitryOlshansky 7y ago
Mostly good stuff but a few comments:
- article doesn’t clarify if it’s on hardware or VMs
- 140 shards per node is certainly on the low side, one can easily scale to 500+ per node (if most shards are small, typically power law distribution)
- more RAM is better, and there is a ratio of disk:ram that you need to keep in mind (30-40 for hot data, 200-300 for warm data)
- heaps beyond 32g can be beneficial but you’d have to go for 64g+, 32-48g is a dead zone
- not a single line about GC tuning (I find default CMS to be quite horrible even in recommended ~31g sizes)
- CPUs are often a bottleneck when using SSD drives
- detaro 7y ago> - heaps beyond 32g can be beneficial but you’d have to go for 64g+, 32-48g is a dead zone I'm curious why that is the case?
- wolf550e 7y agoIf the heap is under 32GB, the JVM can use "compressed oops" with 4 byte "pointers" to 8 byte offsets inside the 32GB heap instead of 8 byte pointers. This makes everything more compact. Since Java doesn't have structs, everything is a pointer and arrays of things are arrays of pointers. So a 48GB heap won't fit a lot more data than a 32GB heap. See https://shipilev.net/jvm/anatomy-quarks/23-compressed-references/ https://shipilev.net/jvm/anatomy-quarks/23-compressed-refere... https://wiki.openjdk.java.net/display/HotSpot/CompressedOops https://wiki.openjdk.java.net/display/HotSpot/CompressedOops
- DmitryOlshansky 7y agoJava compressed object pointers. Shipilev (JVM performance engineer) explains it in great detail here: https://shipilev.net/jvm/anatomy-quarks/23-compressed-references/ https://shipilev.net/jvm/anatomy-quarks/23-compressed-refere... But the short of it: Java uses 32-bit numbers for object references if it can address whole heap with it. Given default 8-byte alignment of object we have 2^32 * 8 = 32g of addressable heap. Once we are beyond that number 64-bit references are used and suddenly we are wasting a lot of heap space (Java is very much objects everywhere). Usually up to around 40% of heap is references.
- tomnipotent 7y ago> Usually up to around 40% of heap is references. This blows my mind. I figured it would be non-trivial, but never expected that much.
- detaro 7y agoThanks, I somehow thought that was at smaller sizes but makes sense of course.
- joking 7y agoSi it’s better to have 2 nodes with 31gb each than one with 64gb?
- DmitryOlshansky 7y agoSee my other reply on big iron vs small instances (as in 31g of heap “small”). In short there are many variables at play, but without any context generally I’d recommend trying many small instances first. Now depending on your hardware, scale and application things may easily change in favor of big iron.
- joking 7y agowell, I was thinking more in using more than one ES instances on the same machine listening on different ports, but your points on the other reply still apply.
- DmitryOlshansky 7y agoI actually did that once - splitting one NUMA machine into 2 ES instances each isolated to its own socket (numactl etc). Sadly at the time there was a big perf problem with indexingdocuments with 20kb+ binary blobs (non-indexed field) and so the change to 2 nodes vs one did have little noticeable effect on the throughput benchmarks I did with our data back then (~2015 ES 1.4) In the end I dropped this split and focused on the other bigger problems. Would love to revisit this exercise with more of free time and better tools.
- DmitryOlshansky 7y agoAnd another note on shards - indexing a shard is a single writer process. If your drive tolerates parallel writes well (=SSD) having multiple primary shards per node helps scale indexing.
- jasontedor 7y agoJason from Elastic here. A shard can handle concurrent writes from clients, there are multiple write threads scaling to the number of logical processors exposed to Elasticsearch, and we have invested in making Lucene accommodate this concurrency[0][1][2]. [0]: http://blog.mikemccandless.com/2011/05/265-indexing-speedup-with-lucenes.html http://blog.mikemccandless.com/2011/05/265-indexing-speedup-... [1]: https://issues.apache.org/jira/browse/LUCENE-3023 https://issues.apache.org/jira/browse/LUCENE-3023 [2]: http://blog.mikemccandless.com/2017/07/lucene-gets-concurrent-deletes-and.html http://blog.mikemccandless.com/2017/07/lucene-gets-concurren... Disclaimer: I am an engineer on the Elasticsearch team; I welcome any and all feedback.
- DmitryOlshansky 7y agoI stand corrected then. I’d certainly need to check my measurements on this one.
- karlney 7y agoHi, do you have real use experience running elasticsearch with 64g+ heap? Is there any articles/benchmark/notes or anything that you would be willing to share? We have considered trying out 64g+ heaps for our cluster but we are concerned about very long gc pauses impacting the search performance.
- wbl 7y agoJava now has pausless GC.
- jjirsa 7y agoIt’s highly concurrent but not pauseless
- StreamBright 7y agoMost production systems are barely on JDK 8. G1GC is the most used GC for high-performance production systems at many large companies (1B+ USD revenue) I have worked for.
- DmitryOlshansky 7y agoI highly recommend trying out Java 11. G1 got quite a few improvements since 8. In particular full GC is parallel since 10: https://openjdk.java.net/jeps/307 https://openjdk.java.net/jeps/307
- StreamBright 7y agoThanks, trying out is not an option for so many things. It is up to vendors to decide which JVM to recommend and I cannot overrule them. As of personal use, I am on 11.
- DmitryOlshansky 7y agoCurrently we are running ES on big iron in production. There are both advantages and disadvantages. On plus side: - easier to manage several big machines than shitload of small instances potentially on top of virtualization solution (debugging overlay network etc. at night is not fun) - nodes are more resilient to big requests (both big bulk indexing and resource-heavy searches). Your risk of hitting sudden OOM is practically zero (although ES does have circuit-breakers that try to prevent processing requests that would cause OOM) - while not having compressed oops is sad, not having multiple copies of JIT and compiled code cache etc. is a plus - using a modern concurrent GC is more sensible on bigger heaps. G1 is actually fine on 64g. Some of us are running Shenandoah on Java 13 in production, I’m looking forward to apply it on my clusters (or get back to testing ZGC). Disadvantages: - NUMA. Try to pick more recent Java as it’s tuned to run better on NUMA machines, still not all of JVM is NUMA-aware - fault domain is larger and getting hot node up to speed (in-sync) after restart can easily take about an hour. Getting it up from empty storage going to take a bloody while - elastic.co folks seem to have recommended settings for lots of small nodes not big ones, so you are on your own to discover proper limits for every setting
- hilbertseries 7y agoI'm kind of surprised this article doesn't mention anything about how many nodes you want in your cluster. Since ES performance starts to degrade once you get past 40 or so nodes .
- DmitryOlshansky 7y agoCertainly a problem with 100-s of nodes. One thing is that GC pauses in big clusters become more problematic as the chance to hit stop-the-world pause on at least one node for every request is getting higher. Going beyond dozens takes a fair amount of tuning and trouble-shooting.
- StreamBright 7y agoYep, I can confirm the GC part. We start with that before touching anything else to get the most out of the system. G1GC is pretty tunable.
- vosper 7y agoWhat kind of performance improvements have you achieved through GC tuning? I'm surprised it's significant enough to start there. Might be something I should look more closely at. Also, have you looked at Shenandoah or ZGC?
- MuffinFlavored 7y agoSerious question: does indexing Logstash/JSON logs really need to take gigabytes of memory + disk and sharding?
- DmitryOlshansky 7y ago(Setting aside the fact that Logstash is JRuby/Java app easily eating said gigabytes of heap) Do JSON logs take gigabytes? If they do for you then yes, gigabytes of disk and memory are pretty much guaranteed. Also things tend to pile up with time (even on a few weeks horizon). In all honesty, I believe a finely crafted native code solution for this problem could achieve x3 less ram usage and x2-3 indexing/search performance. Going beyond that is also possible but will take remarkable engineering. Update: to expand on the last point, c++ solutions are typically closed sourced. Rust and Go both have interesting open-source full text engines: https://github.com/tantivy-search/tantivy https://github.com/tantivy-search/tantivy https://github.com/blevesearch/bleve https://github.com/blevesearch/bleve In near future I totally see someone producing a great open-source distributed search project that is at least on par with today’s ES core feature set.
- geggam 7y agoNot sure if you are aware of this. We ran this at Y! https://vespa.ai/ https://vespa.ai/
- DmitryOlshansky 7y agoSeen that and starred probably half a year ago. Never had the time to dig around and see what it’s like in perf, operations and scalability.
- geggam 7y agoSaw 4 petabyes in it. Cant remember the hosts number exactly but it was several hundred bare metal servers of various sizes Behind Y! Groups
- 7y ago
- holoduke 7y agoIs your list not entirely depending on the usecase? I am using ES for years with over 1 million daily users. It provides simple search funtionality. it runs on a single node with 4gb of memory. For more than 5 years with hardly any issues.
- DmitryOlshansky 7y agoIf the dataset it fits on one node - good. You might still be missing out on reliability (should that node fail or just plain restart). Simple catalogs and/or web pages search usually fit in RAM, so the advice of disk:ram becomes less relevant, also the type of drive would be mostly irrelevant.
- buttheyare 7y agoWhat is your use case and/or dataset size if I may ask?
- paulddraper 7y ago> 140 shards per node is certainly on the low side That seems high, no? Unless you were planning on scaling 20x, it seems you could easily have half the number.
- DmitryOlshansky 7y agoThe total number of shards per node, across all indices. As far as number of shards per index goes there are many considerations some good ones are outlined in the article.
- jillesvangurp 7y agoThe article isn't that good because it mostly just verbatim repeats (some of) the information in the official documentation but sadly mixes it with a lot of things that are simply not correct/misunderstood. Also it omits a lot of stuff that is actually important. The hierarchy breakdown in the article is misleading. Lucene indexes fields, not documents. Understanding this is key. More fields == more files. Segments are per field not per ES index. A lucene index is not the same as an Elasticsearch index. Segments are not immutable but an append only file structure. Lucene creates a new segment every time you create a new writer instance or when the lucene index is committed, which is something that happens every second in ES by default and something you could configure to something higher. So, new segment files are created frequently but not on a per document basis. ES/Lucene indeed constantly merge segment files as an optimization. Force merge is not something you should need to do often and certainly not while it is writing heavily. A good practice with log files is to do this after you roll over your indices. With modern setups, you should be reading up on index life cycle management (ILM) to manage this for you. The notion that ES crashes at ~140 shards is complete bullshit that is based on a misunderstanding of the above. It depends on what's in those shards (i.e. how many fields). Each field has its own sets of files for storing the reverse index, field data, etc. So, how many shards your cluster can handle depends on how your data is structured and how many of them you have. This also means you needs lots of file handles. Understanding how heap memory is used in ES is key and this article does not mention the notion of memory mapped files and even goes as far as to recommend the filesystem cache is not important!! This too is very misguided. The reality is that most index files are memory mapped files (i.e. not stored in the heap) and they only fit in memory if you have enough file cache memory available. Heap memory is used for other things (e.g. the query cache, write buffers, small data-structures with metadata about fields, etc.) and there are a lot of settings to control that which you might want to familiarize yourself with if you are experiencing throughput issues. Per index heap overhead is actually comparatively modest. I've had clusters with 1000+ shards with far less memory. This is not a problem if those shards are smallish. The 32GB memory limit is indeed real if you use compressed pointers (which you need to configure on the JVM, ES does this by default). Far more important is that garbage collect performance tends to suffer with larger heaps because it has more stuff to do. The default heap settings with ES are not great for large heaps. ES recommends having at least (as a minimum) half your RAM available for caching. More is better. Having more disk than filecache means your files don't fit into ram. That can be OK for write heavy setups and might be OK for some querying (depending on which fields you actually query on). But generally having everything fit into memory results in more predictable performance. GC tuning is a bit of a black art and unless you have mastered that, don't even think about messing with the GC settings in ES. It's one of those things where copy pasting some settings somebody else came up with can have all sorts of negative consequences. Most clusters that lose data do so because of GC pauses cause nodes to drop out of the cluster. Mis-configuring this makes that more likely. CPU is important because ES uses CPU and threadpools for a lot of things and is very good at e.g. concurrently writing to multiple segments concurrently. Most of these threadpools configure based on the number of available CPUs and can be controlled via settings that have sane defaults. Also, depending on how you set up your mappings (another thing this article does not talk about) your write performance can be CPU intensive. E.g. geospatial fields involve a bit of number crunching and some more advanced text analysis can also suck up CPU.