How Netflix uses Druid for realtime insights
netflixtechblog.com
netflixtechblog.com
Seriously? I know it's the Medium hustle but someone at Netflix should know better.
Medium isn't worth oodles, though. Looking at Squarespace and Wordpress pricing, for a company's tech blog, it's worth $25-$100 per month. Unless there's value from the Medium brand...
Druid can easily be extended through available 3rd party extensions and you can write your own to implement custom serialisation formats, aggregations, connect to new streaming systems, read directly from whatever cold storage you have etc.
In the Clickhouse model you have to work out a lot more of that stuff yourself though these days it can read from Kafka directly which is useful.
Some things that are important for Clickhouse vs Druid at big scale is the rather large difference in indexing approaches. Clickhouse uses bloom filters and other probabilistic data structures to index large chunks of data, for the most part though actually checking for rows requires a full scan of that chunk to strip false positives.
This is different to Druid which uses full inverted indices for dimension filtering.
The tradeoff is basically Clickhouse is cheaper, especially when scaling out but Druid is faster especially when the cluster is under heavy concurrent query load, like serving analytics dashboards or data exploration interfaces to users.
Clickhouse excels when you want to scan most but not all the data most of the time. Namely reporting or bulk analytics queries that will hit most rows in a block.
I consider both to be excellent databases.
Druid complexity is coming down a bit compared to where it started. These days you need brokers, middlemanagers and historicals - for queries, ingestion and storage respectively.
In the past to do batch ingestion it also required Hadoop but there is now a native parallel batch ingestion system that runs on the middlemanagers as worker tasks that can read from S3/GCS/existing Druid segments.
Druid is by far the more complex but you get a lot for it and with k8s it's not as hard to run/manage as it was in the past.
That being said we are currently working on reducing the number of processes to 4 (from the current 6) for a "standard" setup. The main reason is that at smaller scale there isn't as much of a purpose to having a larger number of processes.
We're also working on removing some of the knobs. Actually, depending on what version you originally looked at, many of them might already be gone.
Also, I don't think they use bloom filters for the index as far as I can tell from the documentation. There is certainly an option to use a bloom filter aggregator on a table for faster counts, but it's not the default. If you're referring to the fact that count () is not precise, there's a exact count function too. This is my speculation, though, and you may be fight.
I will need to check out the materialised views. :)
Let's say you have data with a few dozen dimensions, and want to compute aggregations filtered by any user-supplied union or intersection of dimension values. This is a fairly common use case in analytics dashboards. How do materialized views help with that?
One thing I wanted to add with regard to performance. Druid does indeed get a big boost from the fact that it uses inverted indexes for filtering. It also gets a boost from having a wide variety of approximate algorithms you can use if you want (for things like topN, count distinct, set difference/intersection, quantiles, etc). But straight scan performance has been improving quite a bit recently too.
The biggest change related to straight scan perf is fully vectorizing the query engine, which is partially done as of the latest release (0.17): https://druid.apache.org/docs/latest/querying/query-context..... In benchmarks, the implementation so far has been posting row scan rate improvements in the 2-3x range. I expect we'll be able to round it out and have it work for all queries over the next couple of releases. The multiples involved mean this is quite meaningful if you do a lot of straight scans.
There's plenty of other stuff going on too: our latest release added parallel merging of large result sets. Our next one (0.18) is going to add a new, more efficient hash aggregation engine. That next release is also going to add a JOIN operator -- not perf related, but probably the number one most requested feature.
Vectorized query engine and JOINs sounds awesome.
(We did meet in SF! Beer hall!)
Inverted indexes map distinct values in a column to a list of document ids containing the value. Bitmap indexes map distinct values to an array of booleans the same length as the number of documents, with true for presence and false for absence. Both index types can be highly compressed, of course.
Can you clarify what Druid is using?
How do you account for the possibility that the update only performs badly because it’s different than what users are used to, but would actually be an improvement in the long run?
For something that is effectively binary and can't be changed incrementally the only way is to extend the evaluation window or make a call with limited information.
If you are convinced it's better for the long tun then what's the point in measuring ?
That's assuming that the update contains a UI/UX change. A lot of updates that will roll out won't include that, they'll be fixing or optimising things.
Maybe this is why they finally figured out the auto-playing trailer stuff was a horrible change (I cancelled my account over that). One can only hope.
Even better: sites such as hn should never allow links to sites employing dark patterns.
If people start to care about being exploited the content will move to a more ethical site. Or it might influence medium.
If 1% don't care and share the crap the other 99% still have to suffer for it.
I do agree that medium is a bit of a mess where one is „less privileged“ when signed in.
But what we should be asking for is the companies like Netflix not using medium in a first place.
Netflix's workload would likely exhaust the resources of even a vertically-scaled single node.
Because of that, you don’t really need to scale Materialize in a sharding sense. You can just have a bunch of “the same” Materialize node (i.e. every node just freestanding clone of a template node, with exactly the same sources and matviews) and then hit them with the parts of a map-reduce query launched by, say, Citus—where Citus was thinking it was talking to a bunch of Citus shard nodes each holding a table-shard named X, but was actually talking to Materialize nodes each holding a matview named X. As long as the query sent from the map-reduce job to each node is constrained in its WHERE clause to only the part of the data it expects to get from that node—rather than relying on the node to know what data it has—then the Materialize nodes would each just do the work required to supply that data (including only pulling in the parts of the configured sources required to compute that result.)
That’s just my intuition from how Materialize presents itself as PG-wire-protocol compatible, though; I haven’t tried this myself, and there might be some footguns in the path of anyone really trying to implement it.
And, of course, this is all irrelevant the moment you write a query that needs a pure reduce (e.g. the computation of a current finite-state-machine state over an event-stream source) rather than a map-reduce. Druid/Clickhouse/etc. can probably “scale” those, in at least the Hadoop “move the job around between each serial stage, so each stage has data-locality for the data of that stage” sense; while Materialize would give you no benefit at all in such a job over just querying a plain PG view defined on top of a Foreign Data Wrapper source.
This is absolutely correct!
> You can just have a bunch of “the same” Materialize node (i.e. every node just freestanding clone of a template node, with exactly the same sources and matviews) and then hit them with the parts of a map-reduce query
This should work, but we have been thinking about it/testing it differently internally. In general you should be able to create materialized views on different "shards" that have different `where` conditions, allowing you to control memory that way. This technique does require data that is actually partitionable in this way, same as it must be partitionable in mapreduce.
> this is all irrelevant the moment you write a query that needs a pure reduce
Of course, with materialize's sinks you can spin up a bunch of `materialized`s and connect them for a final reduce after data has gone through e.g. kafka or shared files. Being able to write joins and aggregates across heterogenous sources makes this kind of workload actually pretty pleasant.
For now. We have a pretty good idea of what needs to be done to shed state to disk, and have designed to be able to implement it. We expect it to "just" be a matter of putting in the engineering effort.
It's early days, but I can confirm that it's a bit hard to know out of the gates whether Materialize can handle such workloads. In particular, unless I missed it the query workload isn't really discussed in the post, and that is all that matters.
Druid makes a few compromises that Materialize isn't willing to make. For example AFAIK Druid doesn't support deletions or modifications in "realtime", which means they can track "min/max" style queries much more efficiently, but they lose out on some other features (connecting pipelines of these views where you might have to "retract" a prior min or max). This could certainly let it scale more, but also means that you learn a few months down the road that it doesn't do everything you hoped it would.
No joins in Druid is also reportedly a bit of a pain. Instead, you get to pre-denormalize your data. No need with Materialize (and hey, feel free to push down reductions through the joins, rather than denormalize and then reduce).
Short version, Materialize is definitely a "higher sophistication" data processing play. Whether that works out well remains to be seen! I'll slap up a blog post later today with a worked example.