3 ms·
Yep. There are always going to be constraints about how well a system like clickhouse can support arbitrary joins. Queries in clickhouse are fast because the da
by necubi 3y ago
Yep. There are always going to be constraints about how well a system like clickhouse can support arbitrary joins. Queries in clickhouse are fast because the data is laid out in such a way that it can minimize how much it needs to read.
Part of this is the columnar layout that means it can avoid reading columns that are not involved in the query. However it’s also able to push query predicates into the table scan, using metadata (like bloom filters) that tell it what values are in each chunk of data.
But for joins, you typically end up needing to read all of the data and materialize it in memory.
For realtime joins the best option is to do it in a steaming fashion on ingestion, for example in a system like Flink or Arroyo [0], which I work on.
[0] https://github.com/ArroyoSystems/arroyo https://github.com/ArroyoSystems/arroyo
- closeparen 3y agoSomething I have found pretty annoying is that Flink works great for joining a stream against another stream where messages to be joined are expected to arrive within a few minutes of each other, but there is actually ~no platform solution for joining a small, unchanging or slowly changing table against a stream. We end up needing a service to consume the messages, make RPC calls, and re-emit them with new fields.
- jgraettinger1 3y agoOur (Estuary; I'm CTO) streaming transformation product handles this quite well, actually: https://docs.estuary.dev/concepts/derivations/ https://docs.estuary.dev/concepts/derivations/ Fully managed, UI and config driven. Write SQLite programs using lambdas that are driven by your source events, in order to join data, do streaming transaction processing, and no doubt lots of other things we haven't thought of.
- necubi 3y agoThere are a few options here, but I agree this is a weakness with existing systems. One option in Flink is to load the entire fact table into the pipeline (using the filesystem source or a custom operator) and join against that. This provides good performance, but at the cost of managing additional long-running state in the pipeline (and potentially long startup times). This works pretty well for very small fact tables (stuff like currency conversions, B2B customer data, etc.). The other option is to store the fact table in a database and query it dynamically. Flink SQL has explicit support for this (called "lookup joins") but this requires careful tuning to not overwhelm your database with high-volume streaming traffic (particularly when doing bootstrapping or recovery). Doing these sorts of joins is a huge use case, and definitely something we're trying to improve in Arroyo.
- gunnarmorling 3y agoYou can join against a static (or slowly changing) table in Flink, including AS OF joins [1], i.e. joining the correct version of the records from the static table, as of the stream's event time for instance. You need to keep an eye on state size and potential growth of course. It's a common use case we see in particular for users of change data capture at Decodable (decodable.co). [1] https://nightlies.apache.org/flink/flink-docs-master/docs/dev/table/sql/queries/joins/#event-time-temporal-join https://nightlies.apache.org/flink/flink-docs-master/docs/de...