How Discord Indexes Billions of Messages Using Elasticsearch
blog.discordapp.com
blog.discordapp.com
edit: Put together a logging cluster consisting of 14 nodes, 12 data + 2 indexer/search API, with ~40TB of consumer-grade SSDs. Ran ~14k indices based on log type & timestamp, with a whole raft of custom field configuration to handle aggregations and different tokenizations. In sum, Elasticsearch looks easy to configure and tune but is amazingly hard to do well - but incredibly rewarding.
Isn't this true of everything. The greater the reward the more effort required in most scenarios? I do agree, Elasticsearch can take some trial and error along with experimentation to get right!
- Why only one elasticsearch shard (+replica) per index? You run some major risks of write-time hotspotting here, especially if you ever have to completely rebuild your indices from scratch (pro tip: plan for it now). If you have any basic metric of a discord server's size prior to index creation, you can pre-optimize the number of shards in its index to ensure that writes are distributed across the cluster.
- What happens when one server's index grows too large to fit on one (elasticsearch) shard? Practical shard size limit in ES is ~50GB due to concerns related to reallocation and recovery. You could also theoretically hit the 32GB memory limit.
Also, as another poster mentioned, consider dedicated master nodes. You're going to lose data nodes from time to time, don't let it cost you consistency by bringing down a master with it.
Right now, we're running them under OpenStack on top of our own bare metal with SAS disks. It works well but I have been working on a plan to migrate them to live under Kubernetes like the rest of our infrastructure. I think the answer is to put them in StatefulSets with local hostPath volumes on SSD.
We use CoreOS for our Kube cluster and are huge believers in it so I will definitely be leveraging that for ES. One thing I have done for my Cassandra clusters is to write a custom health check running in a container that monitors Cassandra's health and locks the cluster with locksmithctl to prevent reboots when the database is unhealthy. This will be easy to translate to ES. Correspondingly, I will add a curl call to the ES unit file to disable shard reallocation when the service is shut down during a reboot. This prevents costly reallocation of shards when you know that a node will be back online shortly.
Is it a Medium thing? Does it disallow linking out from a Medium hosted, privately branded blog?
3 small boxes, master eligible, holding no data.
I love the writeup, having a way to reproduce data if your clusters fail is a definite must with ES.
- Having many clusters and assigning messages to a specific cluster seems like an interesting solution.
- I'm curious how they managed to lazily index messages.
- Since only message, channel and server ids are stored in ES, have there been any problems reindexing data after an index fails?
The worst case to an index failure is that the search query is delayed as the index rebuilds itself. We throttle the rate of historical indexing into ES to a safe level so that we're not degrading performance of other components of the system.
Are you using DB triggers to fire the job ?
I'm wondering how long does it take to execute the ES refresh on a search query when the Shard was marked as dirty?
If the search requests are mostly real time, I suspect this is really short, but if the Shard ingest new messages for a while (let's say 50 minutes) and it's marked as dirty, a search query would ask ES to refresh 50 minutes worth of documents before running the actual query.
As it shown to be a problem? Is the refreshing time growing along with the number of documents inserted since the last refresh?
I'm assuming they mean "I can add a new node to the cluster, and NEW SHARDS FROM NEW INDEXES can be distributed amongst the new nodes".... So far as I know elasticsearch can't rebalance shard location or composition automatically based on cluster membership change events...right? I mean, the cluster reroute api will let you manually move a shard, but thats all I know of.
[1] https://www.elastic.co/guide/en/elasticsearch/reference/2.3/...
Fiddling with cluster.routing.allocation.balance.* can really paint you into a corner quickly.
The reviewed tools were: clucene, Apache Lucy, Indri/Lemur, Xapian, and Zettair. Of those, only Xapian and Zettair could be made to return useful search results within time constraints. Sadly, I didn't get Apache Lucy to build, which along with Xapian is the only option still being maintained. clucene could be built, but delivered erratic results (it found "impeachment" but not "indictment" in the US constitution text), and was too difficult to extend into a full stand-alone app (though it apparently is being used in Informix still).
I have only tried it for an hour or two, and can't speak for its search result quality, but it seems promising.
On my laptop, with an index of 10,000 documents each containing 200 random words, it takes about 110ms to perform a fuzzy search. Half of that is the actual search, the other half is the process/initialization overhead.
It's written in C++, and the core features are performance, simplicity(operation and extension of functionality) and a sane, simple API (the codebase is quite small as well, at about 14k lines of code).
Its similar in some regards to Lucene, but does away with most of the exotic features and instead is optimised for rich queries support, fast execution and extensibility.
It should be on GH sometime in April.
Maybe I should reconsider?
I would imagine you could measure some sort of throughput/response time values to tune this instead of relying on cpu usage because having 100% cpu usage doesn't mean app is poorly performing it could mean app is performing optimally with all the resources available to it. Similarly heap_free shouldn't be measured on its own, neither should be tuned on its own, ideally jvm should be left to its own devices to tune these values based on your inputs in form of params like MaxGCPauseMillis, GCTimeRatio ect.
> Search API: An API endpoint that clients can issue search queries to. It needed to do all the permission checks to make sure that clients are only searching messages they actually have access to.
We've had to implement this exact thing for our index, and were never quite sure if we were doing it right. To "prepend" a security clause/condition into an incoming search query felt awkward and hacky.
Do people not hire small business support technicians anymore? I worked for a guy when I was a teen who took support contacts to manage small business networks. He would call me and send me to fix a router when something broke. He couldn't have charged more than $2K per month, basically just to be on-call. That's 24k a year for an on-call ops tech, that's fairly cheap I think.
(Also, "self-healing"? Is this different from hooking up a service request with timeout to an APC with remote power and rebooting the sucker?)
> Of course, this entire search infrastructure would be incomplete without a way to discover clusters and the hosts within them from the application layer.
((If you want something that isn't complicated to use or maintain, why would you rely on technologies which use complicated decentralized algorithms just to find them??))
from:jhgg "Fallout 4" -Witcher
Is that so? If so, you created a parser for this format yourselves and are converting that to ElasticSearch's Query DSL?