Discovering Anomalies in Real-Time with Apache Flink
mux.com
mux.com
But in event processing, unless you can afford yourself to skip events, how do you deal with that sort of thing, especially if the processing needs to keep track of internal state between events?
I read about event-sourcing, which kinda is a solution to that, but add checkpoints and you have pretty much batch processing again.
You can find more details in the docs:
- https://ci.apache.org/projects/flink/flink-docs-release-1.2/...
- https://ci.apache.org/projects/flink/flink-docs-release-1.2/...
My current understanding is that a event-based (in a very broad sense) system are hard to "replay" in the case of a failure (error in data or just a bug), unless you build additional machinery, which decreases robustness. In contrast the task of making a batch processing system perform fast is easy and much better defined.
And its follow-up post: https://www.oreilly.com/ideas/the-world-beyond-batch-streami...
It's fairly rare that I would even attempt to track internal state in an event-processing system (a node will typically emit/attach all necessary information for the next node), but in cases where I do (real-time numerical calculations), we accept a small percentage of error. We'd typically write a health check around our confidence in that calculation and expose that to systems that need to interact with it.
In terms of errors that are systemic rather than due to malformed events, I'd reprocess out of the queue assuming the timing makes sense. For nodes that care about time, we'd write in checks using the event's timestamp as a guard against processing (failing events outside the time range).
There's also an academic paper that was in VLDB's industry track in 2016, "Consistent Regions: Guaranteed Tuple Processing in IBM Streams", http://www.vldb.org/pvldb/vol9/p1341-jacquesSilva.pdf
A high-availability Flink cluster will often use an Apache Zookeeper cluster to elect a leader Job Manager (coordinator) instance. One or more Task Managers (the systems that actually execute the pieces of a Flink application) discover the current Job Manager leader by querying Zookeeper.
Zookeeper tracks the current leader and running jobs. State for those jobs is written periodically to external storage, often an HDFS cluster. Flink automatic-checkpointing will write application state at fixed intervals to storage and, in the event of a failure, automatically resume from the most recent checkpoint. There's also support for manual savepoints which can be used to restore state when submitting a new job or resuming from catastrophic failure.
Flink provides exactly-once guarantees within the context of the Flink application; any side-effects of your application, such as calls to external services or records written to a database, can happen multiple times if you're recovering from a failure.
Can someone explain why Apache are creating projects that compete with each other? Why not focus on one?
tried-and-true statistical methods like probability density functions
Yahoo EGADS anomaly-detection library
Numenta HTM neural-network anomaly-detection library
We ruled out HTM due to AGPL licensing concerns. It's an interesting product, but wasn't a good fit for us at this point in time. EGADS and other basic statistical methods can actually get you pretty far.