3 ms·
There 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 (usi
by necubi 3y ago
There 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.