Differential Dataflow
timelydataflow.github.io
timelydataflow.github.io
A good introduction to the overall topic is incremental [0] an OCaml library by Jane Street and the corresponding blog posts and videos [1]. They even use it for webapps [2].
[0]: https://github.com/janestreet/incremental
[1]: https://www.janestreet.com/tech-talks/seven-implementations-... (Great to understand fundamentals of an incremental computation graph)
[2]: https://www.janestreet.com/tech-talks/intro-to-incr-dom/
In Metamine,
= signifies Equation, a persistent assignment. Any update to the value on the after which any change to the value results in all dependent values being recomputed.
:= signifies a standard one shot assignment.
The runtime behavior is easy enough to implement, less than 100 lines of Pascal. I worried that the code for a language that supported it would be unreadable, like FORTH, only worse. Imagine if Excel spreadsheets were plain text source code.The alternative to all this nonsense is to just throw everything into clickhouse and build materialized views! The drawback is you can't do complex joins, but for 90% of use-cases, clickhouse materialized views work swimmingly.
Here's a slide deck with examples from a recent presentation: https://altinity.com/presentations/introduction-to-high-velo...
ClickHouse sneaks a functional programming model with lambda expressions into SQL. It's not standard SQL but has enormous flexibility.
But what if this is a mobile app, and requests may be delayed by minutes or even hours due to gaps in cellular connectivity? In that case, we need a strategy to handle "late arriving" data.
Flink has several approaches for this, but they're all based around a "high-watermark." This is basically an estimate at time T that all data from time T-D has arrived. Once the watermark has passed the end of our window we consider it closed and can compute the final value for it. (How you compute the watermark is up to you; typically you use a fixed value but this can be made more accurate if you have out-of-band information).
In addition to the default behavior (where late-arriving data is dropped) you can customize this by specifying a trigger that runs when late-arriving date comes, which can be used to e.g., update an external datastore. However at that point the reconciliation is outside of the scope of Flink.
There are various other ways to organize this within flink, i.e., you can keep the windows open indefinitely, update the internal state of your flink job, and serve directly from the flink state, or you can use a periodic trigger to periodically update your downstream from that state. Obviously this will require a larger state size (in disk or in memory depending on your state backend) since you won't be able to close out the window and have to keep all of the data around.
Much more about all this in the Flink docs (https://ci.apache.org/projects/flink/flink-docs-release-1.11...).
Also a lot of the theory behind this comes from the Google Dataflow paper (https://research.google/pubs/pub43864/) which is also just a great read.
If m2 manages m1 and m1 manages p then m2 manages p via m1, but I could not quite understand what this example output should mean. M1 manages m2 and p? M1 is managed by m2 and p. Neither seems correct.
// define a new timely dataflow computation.
timely::execute_from_args(std::env::args(), move |worker| {
Why would something as fundamental as creating a new dataflow need a comment? In other words: why is that not self-evident from the syntax?Not to denigrate the technical accomplishments. And elegance/readability is a very personal thing. But it does seem jarring when the language constructs don't obviously map cleanly to the design intent.
In this case, the command does not create a new dataflow, it creates a new timely dataflow computation by executing the closure on multiple workers. The computation can then spin zero, one, or many dataflows interactively. Dataflows come and go as the computation runs.
The syntax here doesn't otherwise communicate that to me at all.
Enginerrrd said:
> The syntax here doesn't otherwise communicate that to me at all.
Assuming I'm understanding you correctly, that's my point. If the syntax doesn't communicate intent, it suggests a usability limitation.
frankmcsherry said:
> Naming things is hard.
Agreed. Hence point about readability being a personal thing. In this case though, look again at the comment and the line of code. I'd suggest that very few people, on reading the code without the comment, would guess that it was creating a new dataflow computation.
> In this case, the command does not create a new dataflow, it creates a new timely dataflow computation
Good point, thanks for the correction.
Building realtime, low-latency (< 10 sec) dashboards. The type of things where you previously would have had to wait for several hours or a day for ETL pipelines to crunch through a lot of numbers.
We're also fielding interest for streaming ML applications. Ie, moving from batch models to streaming models.
Also worth pointing out that Differential Dataflow has been around for awhile, while Materialize is fairly young. We're still constantly learning about new applications!
In addition, while Materialize does support connecting to other databases, to power "real-time materialized views", a common architecture we are used for is to present a SQL view on top of streaming systems (such as Kafka or Kinesis).
Here's an overview of Materialize that explains the relationship between Timely Dataflow, Differential Dataflow, and Materialize, starting at 23:20 - https://materialize.io/blog-cmudb/