3 ms·
I get not wanting to add yet-another-system to reduce operational complexity but it seems more economical to use a system like Flink to do a time windowed join
by scottcodie 4y ago
I get not wanting to add yet-another-system to reduce operational complexity but it seems more economical to use a system like Flink to do a time windowed join and emit single records to be written to a persistence store. The Flink time window can be sufficiently large to encompass the disparity between ingest and event time without much RAM consumption by using a RocksDB state backend on the operator. Let me know if I miss something, every use case is different :)
- majormajor 4y agoThey don't go into much of the detail of their "event consumers" but it certainly sounds like something Flink can handle, although Flink itself can be yet another operational set of headaches on top of Kafka. It also seems like something simple enough (simple processing-time windows) that they might not need Flink for the consumer. I have a hard time recommending Flink over dumber consumers or a more micro-batch vs true streaming approach unless you're doing something that really needs the long-lived in-worker keyed state and the ability to do things like streaming joins and all. Otherwise the Flink gotchas and nasty surprises can outweigh the ease of which it lets you do what you want to do.
- dikei 4y agoI have written Flink pipelines that solved similar problems to this post: sessionization of time-skew data points from multiple sources with variable and possibly large delays; simple and economical is not how I would describe it :) . If one of your datasource get lag behind, Flink would buffer a huge amount of data waiting for it to catch up due to how watermark work when joining 2 streams, and you would still encounter out-of-memory error even with RocksDB from time to time if your session window get too large. In addition, with our state size frequently reached hundred of GBs, recovering from failure was not exactly fast either.
- DeathArrow 4y agoAs I understood, they wanted to eliminate the MQ layer. If they used Flink, they would still need to keep Kafka and they would be introducing another layer of complexity with Flink. So, instead of simplifying they would make the stack more complex.
- skinnyarms 4y agoThis, especially when some records can be "several megabytes"