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=logical is the setting that enables logical replication. The default (replica) fails when the slot is created: logical decoding requires wal_level >= logical.

  • The postgres user is a superuser, so it has the REPLICATION attribute (a real deployment needs the role of the section above). Connect with host=localhost port=5432 user=postgres password=postgres_password.

  • The image tag is the version of PostgreSQL: 14, 15, 16, 17 and 18 were tested. Change the host port (-p 5433:5432) if 5432 is taken.

  • --rm removes 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#

Record

PostgreSQL

event_type

INSERT / UPDATE / DELETE messages

schema_name / table_name

The relation (Relation message)

primary_key

Columns marked as key in the relation (str or tuple), see primary_keys

record

New row (INSERT, UPDATE) or old row / key (DELETE)

transaction_id

xid of the transaction (Begin message)

global_id

Final LSN of the transaction, e.g. 0/16B3748

event_timestamp

Commit timestamp of the transaction (Unix seconds)

source / position

Name of the slot / LSN of the message (int)

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.execute logs 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_interval seconds, the end of the WAL of the server (wal_end of 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 after wal_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) are EventType.GLOBAL events.

  • 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 with REPLICA IDENTITY FULL (only the old row has them, and the record has the new one): a target must merge an UPDATE into the row it has instead of replacing it. A DELETE (and the old row of a key update) only has the key, unless the table has REPLICA IDENTITY FULL.

  • The values are decoded as UTF-8: use client_encoding=UTF8 in 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.11 and 18.6 (see “Tested Versions”), not with 10 to 13: INSERT, UPDATE, DELETE, jsonb, TOAST columns (with the default replica identity and with FULL), the update of the key, TRUNCATE, REPLICA IDENTITY FULL with primary_keys, the ALTER TABLE detection 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 14 also sends the empty transactions (a Begin and a Commit, they are EventType.GLOBAL events) of the changes of tables that are not published; 15 and 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: Protocol

What 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.

abstractmethod fileno() → int[source]#

File descriptor of the connection, to wait for messages.

__init__(*args, **kwargs)#
class core_cdc.processors.postgres.postgres_wal.TableTruncated(relation: Relation)[source]#

Bases: object

A table was truncated (a Truncate message has one of these for each table).

relation: Relation#
__init__(relation: Relation) → None#
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: IProcessor

It 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", ...)
position: int | None#
transaction_id: str | None#
global_id: str | None#
commit_timestamp: int | None#
stop() → None[source]#

It asks get_events to finish after the current iteration.

get_events() → Iterator[Begin | Commit | Relation | SchemaChange | Insert | Update | Delete | Truncate | Unsupported | TableTruncated][source]#

It returns an iterator with the events to process

get_event_type(event: Any) → EventType[source]#

It returns the event type

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.

name: str#
type_oid: int#
type_modifier: int#
is_key: bool#
__init__(name: str, type_oid: int, type_modifier: int, is_key: bool) → None#
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.

relation_id: int#
schema: str#
table: str#
replica_identity: str#
columns: tuple[Column, ...]#
__init__(relation_id: int, schema: str, table: str, replica_identity: str, columns: tuple[Column, ...]) → None#
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.

previous: Relation#
current: Relation#
__init__(previous: Relation, current: Relation) → None#
class core_cdc.processors.postgres.pgoutput.Begin(final_lsn: int, commit_timestamp: int, xid: int)[source]#

Start of a transaction.

final_lsn: int#
commit_timestamp: int#
xid: int#
__init__(final_lsn: int, commit_timestamp: int, xid: int) → None#
class core_cdc.processors.postgres.pgoutput.Commit(commit_lsn: int, end_lsn: int, commit_timestamp: int)[source]#

End of a transaction.

commit_lsn: int#
end_lsn: int#
commit_timestamp: int#
__init__(commit_lsn: int, end_lsn: int, commit_timestamp: int) → None#
class core_cdc.processors.postgres.pgoutput.Insert(relation_id: int, new: dict[str, Any])[source]#

A row inserted.

relation_id: int#
new: dict[str, Any]#
__init__(relation_id: int, new: dict[str, Any]) → None#
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.

relation_id: int#
new: dict[str, Any]#
old: dict[str, Any] | None#
__init__(relation_id: int, new: dict[str, Any], old: dict[str, Any] | None) → None#
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).

relation_id: int#
old: dict[str, Any]#
__init__(relation_id: int, old: dict[str, Any]) → None#
class core_cdc.processors.postgres.pgoutput.Truncate(relation_ids: tuple[int, ...])[source]#

One or more tables truncated.

relation_ids: tuple[int, ...]#
__init__(relation_ids: tuple[int, ...]) → None#
class core_cdc.processors.postgres.pgoutput.Unsupported(kind: str)[source]#

Any message that is not handled (Origin, Type, Message, streaming, etc.).

kind: str#
__init__(kind: str) → None#
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).

__init__() → None[source]#
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.