Skip to content

Instantly share code, notes, and snippets.

@gfoidl
Last active June 12, 2026 13:48
Show Gist options
  • Select an option

  • Save gfoidl/7ee4643859f28d0825293e945e70d959 to your computer and use it in GitHub Desktop.

Select an option

Save gfoidl/7ee4643859f28d0825293e945e70d959 to your computer and use it in GitHub Desktop.
PostgreSQL replication via WAL
version: '3.8'
services:
db:
container_name: db
image: postgres
restart: unless-stopped
environment:
POSTGRES_USER: root
POSTGRES_PASSWORD: root
POSTGRES_DB: test_db
#PGDATA: /data/postgres
volumes:
- ./data:/data/postgres
ports:
- "5432:5432"
command:
- "postgres"
- "-c"
- "wal_level=logical"
pgadmin:
container_name: pg_admin
image: dpage/pgadmin4
restart: unless-stopped
environment:
PGADMIN_DEFAULT_EMAIL: admin@test.at
PGADMIN_DEFAULT_PASSWORD: root
PGADMIN_CONFIG_CONSOLE_LOG_LEVEL: 30 # warning, cf. https://www.pgadmin.org/docs/pgadmin4/latest/config_py.html#config-py
ports:
- "5050:80"
depends_on:
- db
logging:
driver: none
using System.Text.Json;
using Microsoft.EntityFrameworkCore;
using Npgsql.Replication;
using Npgsql.Replication.PgOutput;
using Npgsql.Replication.PgOutput.Messages;
const string ConnString = "host=10.0.0.20;username=root;password=root;database=test_db";
using CancellationTokenSource cts = new();
await using BloggingContext db = new(ConnString);
await EnsureDbInitialized(db, cts.Token);
Task replicationTask = RunReplication(ConnString, cts.Token);
Blog? blog = db.Blogs.Include(b => b.Posts).FirstOrDefault();
if (blog is null)
{
blog = new()
{
Url = "https://blog.test.at",
Posts = new List<Post>()
{
new Post()
{
Title = $"Post at {DateTimeOffset.Now}",
Content = "Hello world"
}
}
};
db.Blogs.Add(blog);
}
else
{
blog.Posts.Add(new Post()
{
Title = $"Post at {DateTimeOffset.Now}",
Content = "Hello world"
});
}
int rc = db.SaveChanges();
Console.WriteLine($"rc = {rc}, hit any key to continue");
Console.ReadKey();
cts.Cancel();
await replicationTask;
static async Task EnsureDbInitialized(BloggingContext db, CancellationToken cancellationToken)
{
if (await db.Database.EnsureCreatedAsync(cancellationToken))
{
// Cf. https://www.npgsql.org/doc/replication.html
await db.Database.ExecuteSqlRawAsync("create publication posts_pub for table posts;", cancellationToken);
await db.Database.ExecuteSqlRawAsync("select * from pg_create_logical_replication_slot('posts_slot', 'pgoutput');", cancellationToken);
}
}
static async Task RunReplication(string connString, CancellationToken cancellationToken)
{
try
{
await using LogicalReplicationConnection connection = new(connString);
await connection.Open(cancellationToken);
PgOutputReplicationSlot slot = new("posts_slot");
await foreach (PgOutputReplicationMessage message in connection.StartReplication(slot, new PgOutputReplicationOptions("posts_pub", 1), cancellationToken))
{
if (message is InsertMessage insertMessage)
{
await InsertMessageHandler.Handle(insertMessage, cancellationToken);
}
connection.SetReplicationStatus(message.WalEnd);
// Force that status to be updated in the DB. See description of SetReplicationStatus
// for further info.
// For outbox-patterns this avoids that older messages are sent more than once.
// Note: it's still at-least-once and not exactly-once.
await connection.SendStatusUpdate(cancellationToken);
}
}
catch (OperationCanceledException) { }
catch (Exception ex)
{
Console.WriteLine(ex);
}
}
public static class InsertMessageHandler
{
public static async Task Handle(InsertMessage insertMessage, CancellationToken cancellationToken)
{
int colNumber = 0;
Post post = new();
await foreach (ReplicationValue value in insertMessage.NewRow)
{
switch (colNumber)
{
case 0:
{
int? tmp = await TryGetInt(value, cancellationToken);
if (tmp.HasValue)
{
post.PostId = tmp.Value;
}
break;
}
case 1: post.Title = await value.Get<string>(cancellationToken); break;
case 2: post.Content = await value.Get<string>(cancellationToken); break;
case 3:
{
int? tmp = await TryGetInt(value, cancellationToken);
if (tmp.HasValue)
{
post.BlogId = tmp.Value;
}
break;
}
default: break;
}
colNumber++;
static async ValueTask<int?> TryGetInt(ReplicationValue value, CancellationToken cancellationToken)
{
if (value.Kind == TupleDataKind.TextValue && value.GetPostgresType().Name == "integer")
{
string tmp = await value.Get<string>(cancellationToken);
if (int.TryParse(tmp, out int res))
{
return res;
}
}
return null;
}
}
Console.WriteLine(JsonSerializer.Serialize(post, new JsonSerializerOptions { WriteIndented = true }));
}
}
public class Blog
{
public int BlogId { get; set; }
public string Url { get; set; } = null!;
public List<Post> Posts { get; set; } = null!;
}
public class Post
{
public int PostId { get; set; }
public string Title { get; set; } = null!;
public string Content { get; set; } = null!;
public int BlogId { get; set; }
public Blog Blog { get; set; } = null!;
}
public class BloggingContext : DbContext
{
private readonly string _connString;
public BloggingContext(string connString) => _connString = connString;
public DbSet<Blog> Blogs { get; set; }
public DbSet<Post> Posts { get; set; }
protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder)
{
optionsBuilder
//.LogTo(Console.WriteLine)
.EnableSensitiveDataLogging()
.UseNpgsql(_connString)
.UseSnakeCaseNamingConvention();
}
}
@gfoidl

gfoidl commented Apr 15, 2026

Copy link
Copy Markdown
Author

Another example that uses pg_logical_emit_message:

using System.Diagnostics;
using Npgsql;
using Npgsql.Replication;
using Npgsql.Replication.PgOutput;
using Npgsql.Replication.PgOutput.Messages;

const bool UseExplicitTran = true;

/*
Setup for table and replication:

create table if not exists test
(
    id      serial  not null    primary key,
    name    text    not null
);

create publication test_pub for table test;

select pg_create_logical_replication_slot('test_slot', 'pgoutput');
*/

const string ConnString = "host=10.0.0.20;username=root;password=root;database=test_db";

await using LogicalReplicationConnection replicationConnection = new(ConnString);
await replicationConnection.Open();

PgOutputReplicationSlot slot = new("test_slot");

// Having messages: true is a prerequisite for reading messages in Npgsql for pg_logical_emit_message
PgOutputReplicationOptions replicationOptions = new("test_pub", PgOutputProtocolVersion.V4, messages: true);

using CancellationTokenSource cts = new();

Console.WriteLine("Replication starting...");

Task writerTask      = WriteLogicalMessage(UseExplicitTran);
Task replicationTask = Task.Run(async () =>
{
    int id = 0;

    await foreach (PgOutputReplicationMessage message in replicationConnection.StartReplication(slot, replicationOptions, cts.Token))
    {
        string text = $"{++id,3}\tT-ID: {Environment.CurrentManagedThreadId,4}\t{message.GetType().Name,-25} WalStart: {message.WalStart}\tWalEnd: {message.WalEnd}";
        Console.WriteLine(text);

        replicationConnection.SetReplicationStatus(message.WalEnd);
        await replicationConnection.SendStatusUpdate();
    }
}, cts.Token);

await writerTask;
await Task.Delay(250);  // wait a bit, to let the WAL decoder work
cts.Cancel();

try
{
    await replicationTask;
}
catch (OperationCanceledException)
{ }
//-----------------------------------------------------------------------------
static async Task WriteLogicalMessage(bool useExplicitTran, int delayInMillis = 500)
{
    await Task.Delay(delayInMillis);

    /*
        Cf. https://www.postgresql.org/docs/current/functions-admin.html#FUNCTIONS-REPLICATION

        No matter if there's an explicit transation or not, when the transactional parameter is
        set to true, then the output is like:

          1     T-ID:    6      BeginMessage              WalStart: 0/1B92D98   WalEnd: 0/1B92D98
          2     T-ID:    6      LogicalDecodingMessage    WalStart: 0/1B92D98   WalEnd: 0/1B92D98
          3     T-ID:    6      CommitMessage             WalStart: 0/1B92DC8   WalEnd: 0/1B92DC8

        When set to false, the output is like

          1     T-ID:    6      LogicalDecodingMessage    WalStart: 0/1B92E48   WalEnd: 0/1B92E48
    */

    string sql = useExplicitTran
        ? "select pg_logical_emit_message(true, 'heartbeat', 'any message')"
        : "select pg_logical_emit_message(false, 'heartbeat', 'any message')";

    NpgsqlTransaction? tran = null;

    try
    {
        await using NpgsqlConnection conn = new(ConnString);
        await conn.OpenAsync();

        if (useExplicitTran)
        {
            tran = await conn.BeginTransactionAsync();
        }

        await using NpgsqlCommand cmd = new(sql, conn);
        int rc = await cmd.ExecuteNonQueryAsync();

        if (useExplicitTran)
        {
            Debug.Assert(tran is not null);
            await tran.CommitAsync();
        }

        Console.WriteLine($"Emitted message, rc = {rc}");
    }
    catch (Exception ex)
    {
        Console.ForegroundColor = ConsoleColor.Red;
        Console.Error.WriteLine(ex.Message);
        Console.ResetColor();
    }
    finally
    {
        if (tran is not null)
        {
            await tran.DisposeAsync();
        }
    }
}

@gfoidl

gfoidl commented Apr 15, 2026

Copy link
Copy Markdown
Author

View replication / WAL messages from terminal

Setup

create table if not exists test
(
    id      serial  not null    primary key,
    name    text    not null
);

create publication test_pub for table test;

Option A: pgoutput

select * from pg_create_logical_replication_slot('test_slot', 'pgoutput');
pg_recvlogical -d test_db -U root --option=proto_version=4 --option=publication_names=test_pub  --start --slot=test_slot -f -

Then do some DML.

Note

The output is hard to read, as it's kind of binary encoded

Option B: test_decoding

select * from pg_create_logical_replication_slot('test_slot', 'test_decoding');
pg_recvlogical -d test_db -U root -P test_decoding --start --slot=test_slot -f -

outputs something like

BEGIN 763
table public.test: INSERT: id[integer]:1 name[text]:'Egon'
COMMIT 763
BEGIN 764
table public.test: UPDATE: id[integer]:1 name[text]:'foo'
COMMIT 764

# transactional: true
BEGIN 765
message: transactional: 1 prefix: heartbeat, sz: 11 content:any message
COMMIT 765

# transactional: false
message: transactional: 0 prefix: heartbeat, sz: 11 content:any message

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment