Ways to capture changes in Postgres
blog.sequin.io
blog.sequin.io
Here's a quick rundown of how to do it generically https://gist.github.com/slotrans/353952c4f383596e6fe8777db5d... (trades off space efficiency for "being easy").
It's great if you can store immutable data. Really, really great. But you _probably_ have a ton of mutable data in your database and you are _probably_ forgetting a ton of it every day. Stop forgetting things! Use history tables.
cf. https://github.com/matthiasn/talk-transcripts/blob/master/Hi...
Do not use Papertrail or similar application-space history tracking libraries/techniques. They are slow, error-prone, and incapable of capturing any DB changes that bypass your app stack (which you probably have, and should). Worth remembering that _any_ attempt to capture an "updated" timestamp from your app is fundamentally incorrect, because each of your webheads has its own clock. Use the database clock! It's the only one that's correct!
Yes, for consistency you should use the database clock by embedding the calls to `now()` or similar in the query instead of generating it on the client.
But it's not sufficient to use these timestamps for synchronization. The problem is that these timestamps are generated at the start of the transaction, not the end of the transaction when it commits. So if you poll a table and filter for recent timestamps, you'll miss some from transactions that are committing out of order. You can add a fudge factor like querying back an extra few minutes and removing the duplicates, but some transactions will take longer than a few minutes. There's no upper bound to how long a transaction can take in postgresql, and there's a lot of waste in querying too far back. This approach doesn't work if you care about correctness or efficiency.
There's also `statement_timestamp()` per the docs: "returns the start time of the current statement (more specifically, the time of receipt of the latest command message from the client)."
https://www.postgresql.org/docs/current/functions-datetime.h...
This isn't to say any of these methods are the best (or even good) in all cases. Time is tricky, especially if you're trying to do any sort of sequencing.
By default, if you then go on to materialize your collections somewhere else (like Snowflake), you get synchronized tables that follow your source DB as they update.
But! You can also transform or materialize the complete history of your tables for auditing purposes from that same underlying data-lake, without going back to your source DB for another capture / WAL reader.
* Scaling well to multi-TB DBs without pinning the write-ahead log (potentially filling your DB's disk) while the backfill is happening. Instead, our connector constantly reads the WAL and works well in setups like Supabase that have very restrictive WAL sizes (1GB iirc).
* Incremental fault-tolerant backfills that can be stopped and resumed at will.
* Flowing TOAST columns through seamlessly to your materialized destination, without requiring that you resort to REPLICA IDENTITY FULL.
* Being able to offer "precise" captures which are logically consistent in terms of the sequence of create/update/delete events.
The last one becomes really interesting when paired with REPLICA IDENTITY FULL because you can feed the resulting before/after states into an incremental computation (perhaps differential dataflow) for streaming updates of a query.
Our work is based off of the Netflix DBLog paper, which we took and ran with.
[1] https://github.com/estuary/connectors/tree/main/source-postg...
[^1]: https://www.infoq.com/articles/wonders-of-postgres-logical-d...
Yes, marketing and docs are not our strongest suits and we need to do better. To be fair, though, we're also trying not to scare off less technical users who see a bullet list like above and think "well this is clearly not for me". It's a hard balance
I have my own SQLite implementation of a similar pattern (but using columns rather than JSON) which I describe here: https://simonwillison.net/2023/Apr/15/sqlite-history/
For the "capture changes in an audit table" section, I've had good experiences at a previous company with the Temporal Tables pattern. Unlike other major RDBMS vendors, it's not built into Postgres itself, but there's a simple pattern [1] you can leverage with a SQL function.
This allows you to see a table's state as of a specific point in time. Some sample use cases:
- "What was this user's configuration on Aug 12?"
- "How many records were unprocessed at 11:55pm last night?"
- "Show me the diff on feature flags between now and a week ago"
It had been around for decades and over time it had ended up being used for all sorts of things within the company. In fact, it was more or less true that every application and business process within the whole company stored its data within this database.
A key issue we had was that because this database had many different applications that queried it, and there were a huge number of processes and procedures that inserted or updated data within it, sometimes queries would break due to upstream insert/update processes being amended or new ones added that broke application-level invariants -- or when a normal process operated differently when there was bad data.
It was very difficult to work out what had happened because often everything that you looked at was written a decade before you and the employees had long since left the company.
Would it be possible to capture changes from a Postgres database in some kind of DAG in order that you could find out things like:
- What processes are inserting, updating or deleting data and historically how are they behaving? For example, do they operate differently ever?
- How are different applications' querying this data? Are there any statistics about their queries which are generally true? Historically how are these statistics changing?
I don't know if there is prior art here, or what kind of approach might allow a tool like this to be made?
(I've thought of making something like this before but I think this is an area in which you'd want to be a core Postgres engineer to make good choices.)
These approaches work on basically every type of SQL solution that uses WAL/triggers.
For your specific question I have a trigger approach many times in SQL Server but it has a tendency to slow things down if you are logging every query so designing an insertion mechanism that doesn't bog down production isn't perfect, and you might want to perform some sampling.
An approach I've taken is temporal tables w/ Application and UpdatedBy fields. That gives you a permanent record of every change, what application did it, and what user performed the action, at what time, and then lets you query the database as if you were querying it at that point in time. You can add triggers to fail CRUD if those fields are not updated if you want to get really paranoid.
There's a lot of overhead to this in terms of storage, so it's not suitable for high-throughput or cost-constrained transactional systems, but it's something for the toolbox.
You won't get client-level providence data with each change...
However you could hack around that. The logical replication stream can also include informational messages from the "pg_logical_emit_message" function to insert your own metadata from clients. It might be possible to configure your clients to emit their identifier at the beginning of each transaction.
The MermaidJS solves for the DAG and visualizes it, and your browser has enough context to let you CTRL-F for any interesting phrases you put in the label.
https://github.com/pgaudit/pgaudit/blob/master/README.md
https://docs.aws.amazon.com/AmazonRDS/latest/UserGuide/Appen...
[0] https://github.com/arkhipov/temporal_tables [1] https://news.ycombinator.com/item?id=26748096
Event-driven anything is 1000x more complex.
1. Transaction A starts, its before trigger fires, Row 1 has its updated_at timestamp set to 2023-09-22 12:00:01.
2. Transaction B starts a moment later, its before trigger fires, Row 2 has its updated_at timestamp set to 2023-09-22 12:00:02.
3. Transaction B commits successfully.
4. Polling query runs, sees Row 2 as the latest change, and updates its cursor to 2023-09-22 12:00:02.
5. Transaction A then commits successfully.
A simple way to avoid this issue is to not poll close to real-time, as the order is eventually consistent.
Perhaps a more robust suggestion would be to use a sequence? Imagine a new column, `updated_at_idx`, that incremented every time a row was changed.
CREATE OR REPLACE FUNCTION
update_updated_at_function()
RETURNS TRIGGER AS $$
BEGIN
NEW.updated_at = now();
RETURN NEW;
END;
$$ language 'plpgsql';
CREATE TRIGGER
update_updated_at_trigger
BEFORE INSERT OR UPDATE ON
"my_schema"."my_table"
FOR EACH ROW EXECUTE PROCEDURE
update_updated_at_function();
END $$;
Is it possible for two rows to have `updated_at` timestamps that are different from the transaction commit order even if the above function and trigger are used? It's alright if `updated_at` and the commit timestamp are not the same, but the `updated_at` must represent commit order accurate to the millisecond/microsecond.If you're an Elixir & Postgres user, I have a little library for listening to WAL changes using a similar approach:
Would be great if Postgres innovated in this area.
I think that we won’t see traction at the RDBMS “kernel space” until it’s in the SQL standard. There are many valid and complex options to choose from, and there are successful solutions in user space that aren’t overly burdened, performance-wise, from being in user space.
FWIW, the “audit table” approach is the approach that people who study this field gravitate towards. Mainly because it maintains consistent ACIDity in the database, and maintains Postgres as the single point of failure (a trade off vs introducing a proxy/polling job).
I do a lot of work in real-time and streaming analytics. We can do stream processing and maybe some work in the datastore with a materialised view. However, once data has hit your database or data lake you are then effectively back to polling for changes further downstream.
If I want to respond to some situation occuring in my data or update things on a screen without a page refresh then there isn't really a clean solution. Even the solutions in this article feel hacky rather than first class citizens.
Say I want to have a report updating in real time without a page refresh. The go-to approach seems to be to load your data from the database, then stream changes to your GUI through Kafka and a websocket, but then you end up running the weird Lambda architecture with some analytics going through code and some through your database.
There is innovation in this space. KSQL and Kafka Streams can emit changes. Materialize has subscriptions. Clickhouse has live views etc. A lot of these features are new or in preview though and not quite right. Having tried them all, I think they leave far too much work on the part of the developer. I would like a library with the option of [select * from orders with suscribe] and just get a change feed.
I think this is a really important space which hasn't had enough attention.
I know the NoSQL world had some progress in this space. RethinkDB was going in this direction around a decade ago if I recall. I would really like this from a modern relational database or data warehouse though. Polling sucks.
My feeling is that database needs to be built with replication in mind from the get-go to have something like this work well.
Postgres is VERY committed to making sure a replication slot’s consumer doesn’t miss any data. This means that if a consumer stops consuming data from the slot, Postgres will helpfully store all that missed data… right up until the disk fills up and the database falls over. Had this happen during prototyping with two different SaaS DBs, and the only way to get it back up was to file a support ticket. (I can’t remember if metrics warned that the disk was about to fill up or not). Basically, if your replication slot consumer stops reading, that should trigger some kind of alert.
The other reason I don’t use it: the code path to get the initial snapshot of a table is totally different from the code path to read changes. Initializing the read from the replication slot so you never miss any changes is nontrivial.
It’s too bad, because replication is obviously the least hacky solution for change capture.
I use polling, but storing the txid instead of updated_at.
Mind expanding on how you use txid instead of updated_at?
You can configure a size limit after which the slot gets marked as invalid, instead of continuing to retain space. See https://www.postgresql.org/docs/current/runtime-config-repli...
What other behaviour would you like?
> The other reason I don’t use it: the code path to get the initial snapshot of a table is totally different from the code path to read changes.
Hm. For anything dealing with larger data volumes you IME want to handle those things differently (so you can initialize in parallel, initialize from physical backups and similar things). But I can see why it could be useful to optionally stream out existing data out data after slot creation.
> Initializing the read from the replication slot so you never miss any changes is nontrivial.
That part shouldn't be hard - what gave you difficulty?
Unfortunately none of them are perfect as your data scales, each have up and downsides...
I've actually been thinking about turning this idea into a product where you can just point it at your postgres database and select the tables you want to listen to (with filters, like you describe). And have that forwarded to a webhook (with possibility of other protocols like websockets).
I'd love to hear folks thoughts on that (and if it would be something people would pay for). And if anyone might want to partner up on this.
Add your database string, a couple extra migrations, your webhook endpoint - and you're off to the races.
Target customer would be low-code space (think admin dashboards on postgres, webhook integration tools).
https://github.com/supabase/realtime
Spin up a Supabase database and then subscribe to changes with WebSockets.
You can play with it here once you have a db: https://realtime.supabase.com/inspector/new
I remember it being pretty simple (like, run one or two bash commands) to get a source table streamed into a kafka topic, or get a kafka topic streamed into a sink datastore (S3, mysql, cassandra, redshift, etc). Kafka topics can also be filtered/transformed pretty easily.
E.g. in https://engineeringblog.yelp.com/2021/04/powering-messaging-... they run `datapipe datalake add-connection --namespace main --source message_enabledness`, which results in the `message_enabledness` table being streamed into a (daily?) parquet snapshot in S3, registered in AWS Glue.
It is open source but it's more of the "look at how we did this" open source VS the "it would be easy to stick this into your infra and use it" kind of open source :(
Because your API data is in Postgres, you can query it exactly how you like. You're not limited by the API's supported query params, batch sizes, or rate limits. You can use SQL or your ORM.
This also means you don't need to learn all the API's quirks. And every API we support has the same interface. (Which will become more valuable as we add more APIs!)
So, the value is less code and simpler code.
Any other questions, feel free to email: anthony@[domain]
If I remember correctly, the reason is that postgres attempts to de-duplicate messages using a naive quadratic algorithm, causing horrible slowdowns for large query with a FOR EACH ROW trigger.
The way we mitigated that is by creating a temporary table and write the notification payloads. An AFTER * trigger reads distinct rows, groups the data and sends the messages.