4 ms·
Stateful, exaclty once, event processing without the operational capacity to run a proper Flink cluster. This thing needs to be dead simple, pragmatic and cheap
by snerual 5y ago
Stateful, exaclty once, event processing without the operational capacity to run a proper Flink cluster. This thing needs to be dead simple, pragmatic and cheap/simple to operate and update. The only stateful part in our infra at the moment is a PG database.
We are going to start work on this in a weeks, so I'm looking for some insights/shortcuts/existing projects that will make our lives easier.
The goals is to process events from students during exams (max 2500studnets/exam = ~100k-150k events) and generate notifications for teachers. No fancy ML/AI, just logic. Latency of max 1 min.
Our current plan is to let a worker pool lock onto exams (PG lock) and pull new event every few seconds for those exams where (time > last pull & time < now - 10s). All the notifications that are generated are committed together with a serialized state of the statemachine and the ID of the last processed event. Events would just be stored in PG.
This solution is mean to be simple, be implemented in really short timeframe and be a case study for a more "proper & large scale" architecture later on.
Any tips, tricks or past experiences are much appreciated. Also, if you think our current plan sucks, please let me know.
- stingraycharles 5y ago(EDIT: just realised that you specifically mentioned stateful event processing, while what I describe below are two approaches for stateless, exactly-once event processing) Having had a few cracks at this problem, in my opinion using locks is the wrong approach. What you will want is: * split all input data in batches (eg batches of 10k records, or periodic heartbeats every X seconds, etc) * assign each batch a unique identifier * when writing data to the output store, store the batch id along the data; * when retransmitting a batch for whatever reason, reuse the same batch id and overwrite any data in the output store that matches this batch id. Obviously this becomes more tricky when you’re dealing with eg window functions or more complex aggregations. In this situation, I believe that an approach such as “asynchronous barrier snapshotting” works best. Every X seconds, you increment an epoch. While incrementing, you stop ingestion. Then you first tell the output source to create a checkpoint, then the input source to create a checkpoint, and once both have been checkpointed, you can continue streaming data again. Anyway, these are two approaches I’ve used over the years that work well. Explicit locks don’t work well in distributed processing, imho.
- suchow 5y agoI've had good experiences with PQ (https://pypi.org/project/pq/ https://pypi.org/project/pq/). Any event that generates a notification triggers adds an entry to the queue. Worker processes get entries from the queue. The queue is stored as another table in your database whose structure and content is managed by PQ, though you can always read/write to it if you want. PQ handles the concurrency.
- coldilocks 5y agoSounds like a change data capture problem. Consider using Debezium, my team was able to use the standalone java engine to connect to a Postgres DB and stream (within the context of the Java app, not an external kafka stream) insert/update/delete events. You could filter those events and apply your notification and other logic to the filtered events.
- red0point 5y agoI think you could leverage SKIP LOCKED for this - this blog post https://www.2ndquadrant.com/en/blog/what-is-select-skip-locked-for-in-postgresql-9-5/ https://www.2ndquadrant.com/en/blog/what-is-select-skip-lock... explains it nicely.
- analyst74 5y agoMessage queue (i.e. RabbitMQ) sounds like a more natural fit for your problem. What is the peak and avg QPS you need to support? High peak QPS might force you to introduce distributed workers and makes locking impractical. Another consideration is how much do you care about data integrity. Would it be a problem if a few messages are lost? What if a message for processed twice? What if servers lost connection to db for a few seconds? What if a whole server/db goes down?
- benjaminwootton 5y agoDoes Kafka not get you halfway there? It will guarantee exactly once semantics. Use MSK or Confluent cloud if you can use managed services. It’s a more future proof than building this on top of Postgres.
- physicles 5y agoPG is great for this. Should handle ~100 or more events per second without much work (but set up a retention policy, and watch out for tables growing to > ~1M rows, as that will kill you during autovacuum). You can use txid_current_snapshot() and friends to track the last "timestamp". Proper use of locks will help you avoid the complexity associated with long-lived transactions. Exactly-once semantics can be tricky to guarantee if you do it at the wrong layer of abstraction. Sometimes building exactly-once semantics on top of at-least-once semantics is the way to go. Kafka and rabbit MQ are both overkill under 100 events/sec. The extra ops overhead isn't worth it. Besides, with PG it'll be nice to be able to always query a couple tables to completely discern the state of the system.