Last active
June 12, 2026 13:48
-
-
Save gfoidl/7ee4643859f28d0825293e945e70d959 to your computer and use it in GitHub Desktop.
PostgreSQL replication via WAL
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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(); | |
| } | |
| } |
Author
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();
}
}
}
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
Some notes
Publication
In the above SQL the publication is created for all actions, but they can be more specific too, like here where it's filtered to handle only

InsertMessages. But for demo it's kept as is.See e.g. pgAdmin4
Infos about publication / replication