3 ms·
Can you point me to a source for 'routinely stream millions of events per second into druid'? While it is true that druid is great at querying billions of rows
by dataloopio 10y ago
Can you point me to a source for 'routinely stream millions of events per second into druid'?
While it is true that druid is great at querying billions of rows per second it's not very good at ingress. Here is a mailing list discussion for some background.
https://groups.google.com/forum/#!searchin/druid-user/benchmark%7Csort:relevance/druid-user/90BMCxz22Ko/73D8HidLCgAJ https://groups.google.com/forum/#!searchin/druid-user/benchm...
- packetslave 10y agowell, he's is a druid committer and CEO of a company built on top of it, so...
- dataloopio 10y agoNo bias there then
- fangjin 10y agoHere is a source from 2015: http://www.marketwired.com/press-release/metamarkets-clients-analyzing-100-billion-programmatic-events-daily-on-track-surpass-2061596.htm http://www.marketwired.com/press-release/metamarkets-clients... You can also find additional information that folks have been willing to publicly share on scale and use cases here: http://druid.io/druid-powered.html http://druid.io/druid-powered.html
- dataloopio 10y agoThose sources contain literally no technical detail. At 1.1 million metrics per second is that a 40 node druid cluster?
- fangjin 10y agoI think we're using very different terminology here. An event in our world may contain thousands of metrics as part of the same event.
- cheddar 10y agoWhat kind of ingestion numbers are you working with? The thread you link to shows that Druid can ingest ~27.5k events/sec per node, which is roughly 2.376bn events a day per node. While you can claim bias here too, we have multiple clusters ingesting in the high hundreds of thousands of events/second and our largest cluster does close to 2m/s. That's definitely scaled horizontally across multiple nodes. If you are suggesting there is a system out there that can ingest millions of messages a second on a single node, I'd love to hear about it :). edit: Ah, I see from the spreadsheet that you linked that there are systems out there that claim 2.5-3.5m writes per second per node. That's really quite amazing, would be awesome if you could provide the methodology used to collect those numbers. For example, if you are sending in 500 byte events (a rather common size for what we do), if my calculations are correct, you are now sustaining 14 Gbps, which means those benchmarks were done on some beefy hardware. Can you link to a blog post that details the methodology?
- dataloopio 10y agoMost benchmarks are given a colour for reliability and link to repeatability.
- cheddar 10y agoAh, cool, I chased down what you are doing and figured out that you are doing an apples to oranges comparison. As described in your benchmark description: https://gist.github.com/sacreman/b77eb561270e19ca973dd5055270fb28 https://gist.github.com/sacreman/b77eb561270e19ca973dd505527... You are running 200 agents emitting 6000 metrics a piece using Haggar to generate load, which is at https://github.com/dalmatinerdb/haggar https://github.com/dalmatinerdb/haggar The specific thing of interest is how you are generating your data, which looks like you have a single set of dimensions and 6000 metrics dangling off of it. The loop that populates all of the "metrics" are: https://github.com/dalmatinerdb/haggar/blob/master/main.go#L51-L56 https://github.com/dalmatinerdb/haggar/blob/master/main.go#L... And the thing that actually populates the bytes are at: https://github.com/dalmatinerdb/haggar/blob/master/util.go#L21-L37 https://github.com/dalmatinerdb/haggar/blob/master/util.go#L... So, if we take this to an apples-to-apples comparison, you have 200 agents sending a single event every second with 6000 metrics in it. That means that you are successfully ingesting 200 events per second in the way that we would measure event ingestion for Druid. Note, also, that the thread you link to is ingesting 17 independent dimensions with each and every event that flows in. From the Daltaminer docs, it looks like you put all dimension data into postgres and you don't expect any large-scale deployment to ever need more than a single postgres node: https://gist.github.com/sacreman/9015bf466b4fa2a654486cd79b777e64 https://gist.github.com/sacreman/9015bf466b4fa2a654486cd79b7... Look under "Setup Postgres". We routinely have billions of unique combinations of dimension values per day flowing into our system. Delegating the finding of the right keys to a relational database for such operations is going to be very cost-prohibitive, not to mention, you are going to have to materialize hundreds of millions of keys in order to do a simple aggregate over the day. So, I guess this is just another case where you should never trust benchmarks that you didn't do yourself or that don't follow a standard pattern like TPC-H. It's too easy for the same words to be used with different meanings.