Repository navigation
Subscriptions sometimes skips events #222
Description
Activity
Have you confirmed that the handler hasn't fired, like with logs or something?
Yes (sorta), I am logging every MongoDB call and the one to upsert the offending documents was never called
However, digging deeper into the logs I noticed my container was getting OOMKilled which might explain it, however that possibly raises another issue where the checkpoint is updated before the handlers have run?
No it can't be. The checkpoint commit is downstream from the projector. Are you sure you are doing
SaveAsyncorExecuteAsyncand it's not being delayed in any way by the library you use?This is a part of the handler that doesn't always run:
public class RegistrationProjections : EventHandler { public RegistrationProjections() { On<V1.RegistrationReceived>( async ctx => await new Registration { ID = ctx.Message.Id.ToString(), [... Other properties ...] }.SaveExceptAsync(x => new {x.Barcodes}) ); } }Barcodes are being updated in another event, which is why they are ignored here.
The handler is being registered like this:
services.AddSubscription<PostgresAllStreamSubscription, PostgresAllStreamSubscriptionOptions>( "RegistrationProjections", builder => builder.Configure(x => x.MaxPageSize = 256).AddEventHandler<RegistrationProjections>() );I tried lowering the
MaxPageSizeto 256 see if it had any effect, but that doesn't seem to have done anythingWhat if you try using Eventuous MongoDB tools for that upsert? Or the MongoDB driver native API? I just want to remove the possibility that it's a third-party dependency causing the issue.
I understand, however that would require quite the refactor in my handlers.
I will try to add some more debug logging and try to dig deeper into it, and maybe create a simple handler using the Eventuous MongoDB tools to see if I can replicate it there.
I tried putting some more logging into the handler:
public class RegistrationProjections : EventHandler { public RegistrationProjections(ILogger<RegistrationProjections> logger) { On<V1.RegistrationReceived>( async ctx => { logger.LogDebug( "Running projection for {EventType} with ID {RegistrationID}", typeof(V1.RegistrationReceived).FullName, ctx.Message.Id ); await new Registration { ID = ctx.Message.Id.ToString(), [... Other properties ...] }.SaveExceptAsync(x => new {x.Barcodes}); } ); } }Most of the events log it fine, but some are skipped. The event streams does exist in the database and this time there were no OOM kills of the container
Can you add global position to the logs? Also, how do you produce these migration events? Do you commit from multiple processes/threads, or is it a linear fetch-produce?
Here is an excerpt from the Eventuous debug logs - the
Positionproperty is the one highlighted in blue:- Dont mind that the event handler names etc are different from above - the above is merely a sample whereas the below is from real code

Note that positions 246197, 246201 & 246202 are missing even though they exist in the database:
(The ID's are off by 1 because of #163)I have 1 service that discovers what data should be migrated which pushes a bunch of messages to another service using MassTransit that then uses the command service to apply the events and then a third service that generates the projections - the logs are from the third service. So the commits are multi-threaded from multiple containers (running the same code) but the projections are guaranteed to only be run by 1 instance at a time
I see that the global position appears out of order in the logs, do you partition the subscription?
For missing events, my suspicion is that it's the sequence issue. In case of frequent concurrent writes, in once service the sequence might get allocated before and commit later than in the other service. As the result, you can events with the lower sequence number committed after events with the higher sequence number. It results in the subscription receiving the events with higher number first and the next call will say "more than the one I have", and you have skipped events.
I know that only in theory as I am not an expert in Postgres, but I have read about it somewhere. Will try to find more about it, you can spend some time googling too...
Unexpected results might be obtained if a cache setting greater than one is used for a sequence object that will be used concurrently by multiple sessions. Each session will allocate and cache successive sequence values during one access to the sequence object and increase the sequence object's last_value accordingly. Then, the next cache-1 uses of nextval within that session simply return the preallocated values without touching the sequence object. So, any numbers allocated but not used within a session will be lost when that session ends, resulting in "holes" in the sequence.
Furthermore, although multiple sessions are guaranteed to allocate distinct sequence values, the values might be generated out of sequence when all the sessions are considered. For example, with a cache setting of 10, session A might reserve values 1..10 and return nextval=1, then session B might reserve values 11..20 and return nextval=11 before session A has generated nextval=2. Thus, with a cache setting of one it is safe to assume that nextval values are generated sequentially; with a cache setting greater than one you should only assume that the nextval values are all distinct, not that they are generated purely sequentially. Also, last_value will reflect the latest value reserved by any session, whether or not it has yet been returned by nextval.
Basically, what the docs say is to use cache setting of one to guarantee order in the sequence.
As we don't use sequences explicitly, I am wondering what are the sequence settings for the unique auto incremented id...
do you partition the subscription?
Partitioning is not enabled
--
A quick search around suggests that the cache values are only really a thing for a sequence, not for the identity. I will try searching deeper
Ok, looks like I am spamming here, but still.
Basically, what I am trying to say is that the issue is not that the subscription skips events. It's most probably caused by events with higher global position being committed to the database before events with lower values in global position. I am not exactly sure how to solve it.
I thought of the following:
- Don't use autogenerated id
- Query the max id value after the transaction is opened
- Assign incrementing values to global position when inserting
- Commit the transaction
- It might fail in case of optimistic concurrency, then retry
It will probably slow down the appends, but should solve the issue.
10 remaining items
Do you have any suggestions on how to implement a failing test for this issue?
I was able to reproduce the issue by doing the following:
- adding a random delay in AppendEvents.sql before the return statement (to simulate different transactions durations)
- writing a (very rudimentary) test that appends messages concurrently while running a subscription with polling interval set to 0
- checking for gaps in the consume context global position received by the subscription handler
Without this delay, the test fails 1 out of 20 times. With the delay, it fails every time.
Would it be valid to add this random delay, enabled by a new optional parameter, to the AppendEvents function just for testing purposes?
I guess using an optional parameter makes sense. It'd be good to have it 100% reproducible. However, I am not sure what the solution is. One thing I have in mind is to stop using the sequence and instead use the max log position for each append, and fail on uniqueness in case of conflict. However, it effectively means locking the log for concurrent appends.
Another option (speculating here) would be to make sure that the previous log position (current append first log position minus one) is available in the table, and wait until it's there. It would ease the lock as it would still allow concurrent appends but they will be effectively sequenced at the end to make sure that the log position is constantly increasing.
I cleaned up the test a bit: gschuager@8fb8b4d
I also experimented with the non-autogenerated global position approach (gschuager@7ef9d24). It looks like it can work, though it does require retrying on concurrent appends. Still, I believe it's preferable to accept some performance degradation over risking incorrect behavior: subscriptions that skip events aren't reliable and would be a showstopper.
I'm not sure yet where the retry logic should live; do you have any preference?Regarding your suggestion to check and wait for consecutive global_position values: we should keep in mind that PostgreSQL sequences can have gaps (e.g., due to failed transactions), so that might not be a viable strategy.
I am also looking into doing something like what Oskar described in its blog post (using postgres transaction ids to track subscriptions offsets), but that requires some more work.
Upon reflection, I realized that Oskar approach would work if we were to append a single event per transaction: we would simply use pg_current_xact_id() as the global_position
However, since we need to insert multiple events within a single transaction (multiple events with the same tx id), we cannot use just the tx id as he proposes to store subscriptions checkpoints ... the only thing that comes to my mind would be to use the pair [tx_id, event index within the tx] as the checkpoint, but that would require core changes.
Is this still an issue? I see #424 has been merged but not yet released.
I'm happy to support with further testing or further investigation too.
I included test cases for the scenarios that I was able to come up with in the PR, but for sure it could use additional testing. Maybe the issue could be resolved until new evidence comes up.
The preview version should be available on MyGet. Maybe you can test it and see if it works better.
Reacted by Matheus AntunesOkay so I did some manual testing for now on version 0.15.3-alpha.0.30:
I enabled the random sleep implemented by @gschuager (also increased it) in
append_events.sqland implemented anEventHandlerthat stores the event on a flat table.
Then I sent 50+ commands in parallel.Even after thousands of events processed, the stream count and the projection table count are still the same. Also the subscription checkpoint (also postgres) points to the latest processed event.
My assumption here was that given the delay and heavy load, the tracking may miss an event and it's just not the case 🙂.
Ordering is slightly different in both tables (messages and projection) but that's expected given the load.
SELECT 'events_count' AS name, count(*) AS count FROM streams UNION ALL SELECT 'projection_records' AS name, count(*) FROM public.projection UNION ALL SELECT 'messages' AS name, count(*) FROM messages UNION ALL SELECT 'checkpoint' AS name, position FROM checkpoints WHERE id = 'another-sub' UNION ALL SELECT 'global' AS name, MAX(global_position) from messages;
produces
Do you see any other scenario that could be tested?
Reacted by Germán SchuagerI thought this was fixed last year, anyone want to confirm?
I have a funny problem with my setup that is running with the Postgresql EventStore (using npgsql 7) and MongoDB EventHandlers (using MongoDB.Entities) that are really simple - basically only upserting data without any weird data modification along the way.
However, sometimes when loading a bunch of data (fx from importing data from a legacy application) the checkpoint store has passed through a specific event but the EventHandler has never been run. If I then reset the checkpoint to a lower value and restart the application the EventHandlers run fine and everything is up to date.
Any ideas on how to debug this further is much appreciated - I am working on a PoC sample to see if it can be replicated outside our environment