PostgreSQL#
PostgresWalProcessor reads the changes of a PostgreSQL logical replication slot that
uses the built-in pgoutput plugin (PostgreSQL 10 or later), so nothing has to be installed
in the server (it works with the managed services that offer logical replication too).
The processor does not import any driver: it works with the cursor of a logical replication
connection of psycopg2 (psycopg2.extras.LogicalReplicationConnection). psycopg3 does
not have an API for the replication protocol.
pip install '.[postgres]'
Server requirements#
-- postgresql.conf (restart required): wal_level = logical
CREATE ROLE cdc WITH LOGIN REPLICATION PASSWORD '...';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO cdc;
CREATE PUBLICATION cdc_publication FOR ALL TABLES;
SELECT pg_create_logical_replication_slot('cdc_slot', 'pgoutput');
Local server with Docker#
To try the processor (or to run its functional test) a server with logical replication takes one command:
docker run -d --rm --name pg -p 5432:5432 \
-e POSTGRES_PASSWORD=postgres_password \
postgres:16 -c wal_level=logical
# It must print "logical":
docker exec pg psql -U postgres -tAc "show wal_level"
-c wal_level=logicalis the setting that enables logical replication. The default (replica) fails when the slot is created:logical decoding requires wal_level >= logical.The
postgresuser is a superuser, so it has theREPLICATIONattribute (a real deployment needs the role of the section above). Connect withhost=localhost port=5432 user=postgres password=postgres_password.The image tag is the version of PostgreSQL:
14,15,16,17and18were tested. Change the host port (-p 5433:5432) if5432is taken.--rmremoves the container when it stops (docker stop pg); its data is lost, and the slots with it. In a server that stays, drop the slots that are not used (SELECT pg_drop_replication_slot('cdc_slot')).
The functional test (tests/functional/check_postgres_wal.py) creates and drops its own
table, publication and slot, it is not part of quick_test.sh:
pip install -e ".[postgres]"
python -m pytest -s tests/functional/check_postgres_wal.py
# Another port: PORT_TEST_POSTGRES=5433 python -m pytest -s ...
It reads HOST_TEST_POSTGRES, PORT_TEST_POSTGRES, DATABASE_TEST_POSTGRES,
USER_TEST_POSTGRES and PASSWORD_TEST_POSTGRES (the defaults match the command above).
Usage#
import logging
import psycopg2
from psycopg2.extras import LogicalReplicationConnection
from core_cdc.processors.postgres import PostgresWalProcessor
connection = psycopg2.connect(
# The text is decoded as UTF-8: the server must send it in that encoding.
"dbname=tests user=cdc host=localhost client_encoding=UTF8",
connection_factory=LogicalReplicationConnection,
)
cursor = connection.cursor()
cursor.start_replication(
slot_name="cdc_slot",
decode=False, # The processor decodes the binary messages of pgoutput.
options={"proto_version": "1", "publication_names": "cdc_publication"},
)
processor = PostgresWalProcessor(
stream=cursor,
slot_name="cdc_slot",
service="my-service",
targets=[...],
logger=logging.getLogger(),
# Only for the tables with REPLICA IDENTITY FULL (see below):
primary_keys={"public.accounts": "id"},
)
processor.execute() # It runs until `processor.stop()` is called.
psycopg2 2.8.3 or later is needed (send_feedback(force=True)).
How it maps to Record#
|
PostgreSQL |
|---|---|
|
|
|
The relation ( |
|
Columns marked as key in the relation ( |
|
New row (INSERT, UPDATE) or old row / key (DELETE) |
|
|
|
Final LSN of the transaction, e.g. |
|
Commit timestamp of the transaction (Unix seconds) |
|
Name of the slot / LSN of the message ( |
Values are decoded from the text format of PostgreSQL: boolean, integers, floats,
numeric (Decimal), json / jsonb (parsed) and bytea (bytes); any other
type (dates, timestamps, UUIDs, arrays, etc.) is kept as the str PostgreSQL sends.
The primary key. The server only says which columns identify a row, and with
REPLICA IDENTITY FULL that is every column, so it is not the primary key: for those tables
primary_key is () (and a warning is logged, once) unless the key is given in
primary_keys={"schema.table": "id"} (or a tuple of columns), that is also used for the
tables that have no key.
Updates of the key. UPDATE t SET id = 2 WHERE id = 1 gives two records, in this
order: a DELETE with the old row (only the key, the other columns are None, unless the
table has REPLICA IDENTITY FULL) and the UPDATE with the new one. A target that upserts
by key would keep the row of the old key otherwise. (MySQL only sends the new row.)
Confirming the position#
Unlike the BinLog of MySQL or the resume token of MongoDB, the position is kept by the
server in the slot, so the processor has no method to persist it. It confirms to the server
(send_feedback(flush_lsn=..., force=True), psycopg2 only stores the position otherwise)
the end LSN of a transaction, when the consumer is done with its
Commit. Like the other processors do with their positions, it does not know if the targets saved the records (IProcessor.executelogs the errors of the targets and goes on), so a target that failed has to keep the records itself or stop the process;while there are no changes, every
keepalive_intervalseconds, the end of the WAL of the server (wal_endof the cursor, but never inside a transaction): a publication of a quiet table in a busy database would keep all the WAL of the server otherwise. The server closes the connection afterwal_sender_timeout(60 seconds by default).
A message that can not be decoded raises the error and stops the process: skipping it would
lose the change, as the next Commit confirms its position too. The server sends everything
after the last confirmed position again when the process is started. An unconfirmed slot keeps
the WAL in the server: drop the slots that are not used anymore.
Schema changes (DDL)#
PostgreSQL does not write DDL statements in the WAL. The server sends a Relation message
(the definition of the table) before the first change of a table after it was altered, and the
processor compares it with the one it has: when the columns of a table that was already known
changed (or its replica identity) it emits a DdlEvent with DdlAction.ALTER_TABLE,
statement=None (as MongoDB) and the old and new definitions in raw (a SchemaChange),
so the targets must build their own DDL. CREATE and DROP are not detected, nor an
ALTER done while the processor was not running (the definitions are kept in memory).
A TRUNCATE is a DdlEvent with DdlAction.OTHER and statement=None for each
table it truncates (raw is a TableTruncated), so the targets that do not handle it skip
it: a target that keeps the rows has to be told to delete them.
Limitations#
Protocol version 1 only: the transactions are sent when they are committed (no streaming of in-progress transactions, no two phase commit).
The logical decoding messages (
pg_logical_emit_message) areEventType.GLOBALevents.The columns that were not changed and are stored out of line (TOAST, values over ~2 kB) are not sent in the new row of an
UPDATE, so they are not in the record, and that is also true withREPLICA IDENTITY FULL(only the old row has them, and the record has the new one): a target must merge anUPDATEinto the row it has instead of replacing it. ADELETE(and the old row of a key update) only has the key, unless the table hasREPLICA IDENTITY FULL.The values are decoded as UTF-8: use
client_encoding=UTF8in the connection (the server converts the text to the encoding of the connection).It was tested live with PostgreSQL
14.24,15.17,16.13,17.11and18.6(see “Tested Versions”), not with10to13: INSERT, UPDATE, DELETE,jsonb, TOAST columns (with the default replica identity and withFULL), the update of the key,TRUNCATE,REPLICA IDENTITY FULLwithprimary_keys, theALTER TABLEdetection and the confirmation of the position (of a commit and while idle). What is not covered live is covered by the unit and integration tests, that build the messages of the protocol by hand.PostgreSQL
14also sends the empty transactions (aBeginand aCommit, they areEventType.GLOBALevents) of the changes of tables that are not published;15and later do not.
Further reading#
API#
PostgreSQL logical replication (WAL) processor for Change Data Capture services.
- class core_cdc.processors.postgres.postgres_wal.ReplicationCursorProtocol(*args, **kwargs)[source]#
Bases:
ProtocolWhat the processor needs from the cursor of a logical replication connection (psycopg2.extras.ReplicationCursor), so no driver is imported here.
- abstractmethod read_message() Any | None[source]#
A message (with payload and data_start) or None if there is none yet.
- abstractmethod send_feedback(write_lsn: int = 0, flush_lsn: int = 0, apply_lsn: int = 0, reply: bool = False, force: bool = False) None[source]#
It tells the server the position that was processed. psycopg2 only stores the values (the message is sent when its status_interval is over) unless force.
- __init__(*args, **kwargs)#
- class core_cdc.processors.postgres.postgres_wal.TableTruncated(relation: Relation)[source]#
Bases:
objectA table was truncated (a Truncate message has one of these for each table).
- class core_cdc.processors.postgres.postgres_wal.PostgresWalProcessor(stream: ReplicationCursorProtocol, slot_name: str, keepalive_interval: float = 10.0, primary_keys: dict[str, str | tuple[str, ...]] | None = None, **kwargs)[source]#
Bases:
IProcessorIt processes the changes of a PostgreSQL logical replication slot that uses the pgoutput plugin.
The position is kept by the server in the slot, so there is nothing to persist: it is confirmed (send_feedback) when the consumer of get_events is done with the Commit of a transaction (IProcessor.execute logs the errors of the targets and goes on, as the other processors do with their positions), and while there are no changes the end of the WAL of the server is confirmed too (only outside a transaction) so the slot moves even when the publication is quiet. An unconfirmed slot keeps the WAL in the server.
A message that can not be decoded stops the stream (the error is raised): the next Commit would confirm its position and the change would be lost.
More information: https://www.postgresql.org/docs/current/logical-replication.html
- __init__(stream: ReplicationCursorProtocol, slot_name: str, keepalive_interval: float = 10.0, primary_keys: dict[str, str | tuple[str, ...]] | None = None, **kwargs) None[source]#
- Parameters:
stream – Cursor of a logical replication connection, with the replication already started (start_replication(…, decode=False)).
slot_name – Name of the slot, it is the source of the records.
keepalive_interval – Seconds to wait for a message before confirming the last processed position to the server (and how often it is confirmed while there are no changes). It must be lower than wal_sender_timeout of the server.
primary_keys – The primary key of the tables, by
"schema.table". The server only says which columns identify a row, and with REPLICA IDENTITY FULL that is every column: the key of such tables (and of the ones without a key that were configured here) is what is given here, otherwise it is().
The server must send the text in the encoding of the connection, use UTF-8 (
client_encoding=UTF8): the values are decoded as UTF-8.Example:
cursor = connection.cursor() # connection_factory=LogicalReplicationConnection cursor.start_replication( slot_name="cdc_slot", decode=False, options={"proto_version": "1", "publication_names": "cdc_publication"}, ) PostgresWalProcessor(stream=cursor, slot_name="cdc_slot", ...)
- get_events() Iterator[Begin | Commit | Relation | SchemaChange | Insert | Update | Delete | Truncate | Unsupported | TableTruncated][source]#
It returns an iterator with the events to process
- get_ddl_event(event: SchemaChange | TableTruncated) DdlEvent[source]#
It converts the change of the columns of a table (DdlAction.ALTER_TABLE) or the truncation of a table (DdlAction.OTHER) into a DdlEvent. PostgreSQL does not write the statements in the WAL, so statement is None and the targets must build their own DDL from the action, schema and table (the event is in raw, a SchemaChange has the previous and current definitions of the table).
- process_event(event: Any) None[source]#
It should be implemented if another event (apart from DDL and DML) must be processed…
- process_dml_event(event: Any) list[Record][source]#
It converts a change into records. An update that changes the primary key has two: the delete of the old key and the update (a target that upserts by key would keep the row of the old key otherwise). Any other event (if EventType.GLOBAL is streamed IProcessor.execute sends them all here) is processed with process_event.
Decoder of the pgoutput logical replication messages (protocol version 1) of PostgreSQL.
https://www.postgresql.org/docs/current/protocol-logicalrep-message-formats.html
- class core_cdc.processors.postgres.pgoutput.Column(name: str, type_oid: int, type_modifier: int, is_key: bool)[source]#
A column of a relation, as it is described by a Relation message.
- class core_cdc.processors.postgres.pgoutput.Relation(relation_id: int, schema: str, table: str, replica_identity: str, columns: tuple[Column, ...])[source]#
Definition of a table, PostgreSQL sends it before the first change that uses it.
- class core_cdc.processors.postgres.pgoutput.SchemaChange(previous: Relation, current: Relation)[source]#
A table that was already known is described differently (its columns or its replica identity changed). The decoder returns it instead of the Relation message.
- class core_cdc.processors.postgres.pgoutput.Begin(final_lsn: int, commit_timestamp: int, xid: int)[source]#
Start of a transaction.
- class core_cdc.processors.postgres.pgoutput.Commit(commit_lsn: int, end_lsn: int, commit_timestamp: int)[source]#
End of a transaction.
- class core_cdc.processors.postgres.pgoutput.Insert(relation_id: int, new: dict[str, Any])[source]#
A row inserted.
- class core_cdc.processors.postgres.pgoutput.Update(relation_id: int, new: dict[str, Any], old: dict[str, Any] | None)[source]#
A row updated. The unchanged TOAST columns are not in new.
- class core_cdc.processors.postgres.pgoutput.Delete(relation_id: int, old: dict[str, Any])[source]#
A row deleted: the key or the full old row (see the replica identity).
- class core_cdc.processors.postgres.pgoutput.Truncate(relation_ids: tuple[int, ...])[source]#
One or more tables truncated.
- class core_cdc.processors.postgres.pgoutput.Unsupported(kind: str)[source]#
Any message that is not handled (Origin, Type, Message, streaming, etc.).
- core_cdc.processors.postgres.pgoutput.format_lsn(lsn: int) str[source]#
It formats a WAL position as PostgreSQL does (0x16B3748 -> ‘0/16B3748’).
- class core_cdc.processors.postgres.pgoutput.PgOutputDecoder[source]#
It decodes the payloads of the replication stream. It remembers the Relation messages, the only state it has, to name the columns of the tuples and to notice when a table changes (see SchemaChange).
- get_relation(relation_id: int) Relation[source]#
The last definition received for a table.
- Raises:
ValueError – If the relation was not received.
- decode(payload: bytes) Begin | Commit | Relation | SchemaChange | Insert | Update | Delete | Truncate | Unsupported[source]#
It decodes a pgoutput (version 1) payload into a typed message.
- Raises:
ValueError – If the payload is empty or incomplete, or if it has a change of a relation that was not received before.