WalFlow.Sources.Postgres 0.1.0-alpha.alpha.20260708090623

This is a prerelease version of WalFlow.Sources.Postgres.
dotnet add package WalFlow.Sources.Postgres --version 0.1.0-alpha.alpha.20260708090623
                    
NuGet\Install-Package WalFlow.Sources.Postgres -Version 0.1.0-alpha.alpha.20260708090623
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="WalFlow.Sources.Postgres" Version="0.1.0-alpha.alpha.20260708090623" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="WalFlow.Sources.Postgres" Version="0.1.0-alpha.alpha.20260708090623" />
                    
Directory.Packages.props
<PackageReference Include="WalFlow.Sources.Postgres" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add WalFlow.Sources.Postgres --version 0.1.0-alpha.alpha.20260708090623
                    
#r "nuget: WalFlow.Sources.Postgres, 0.1.0-alpha.alpha.20260708090623"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package WalFlow.Sources.Postgres@0.1.0-alpha.alpha.20260708090623
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=WalFlow.Sources.Postgres&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Addin
#tool nuget:?package=WalFlow.Sources.Postgres&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Tool

WalFlow.Sources.Postgres

PostgreSQL Change Data Capture source for WalFlow using logical replication (WAL streaming).

Features

  • Logical Replication - Captures changes via PostgreSQL WAL streaming
  • Multiple Snapshot Strategies
    • TransactionalExported - Consistent point-in-time snapshot with transaction isolation
    • Incremental - Debezium-style chunked snapshot with live WAL merge (non-blocking)
    • BestEffort - Simple snapshot without blocking writes
  • Automatic Publication Management - Creates and manages PostgreSQL publications
  • Replication Slot Lifecycle Management - Configurable create/require/recreate and drop-on-stop behavior
  • Graceful Shutdown - SIGTERM/SIGINT handling with proper checkpoint and flush
  • Compound Primary Key Support - Full support for multi-column keys and keyless tables
  • Keep-Alive Protocol - Proper WAL sender acknowledgment

Installation

dotnet add package WalFlow.Sources.Postgres

PostgreSQL Setup

-- Enable logical replication
ALTER SYSTEM SET wal_level = 'logical';
-- Restart PostgreSQL

-- Create a replication user (optional but recommended)
CREATE ROLE cdc_user WITH REPLICATION LOGIN PASSWORD 'password';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO cdc_user;

Usage

services.Configure<CdcProcessorOptions>(o =>
{
    o.ConnectionString = "Host=localhost;Username=postgres;Password=***;Database=mydb";
    o.SlotName = "cdc_slot";
    o.PublicationName = "cdc_publication";
    
    o.SnapshotMode = SnapshotMode.Initial;
    o.SnapshotStrategy = SnapshotStrategy.Incremental;

    // Replication slot lifecycle (Debezium-style defaults)
    o.ReplicationSlot.ManagementMode = ReplicationSlotManagementMode.EnsureExists;
    o.ReplicationSlot.CleanupMode = ReplicationSlotCleanupMode.Keep;
    o.ReplicationSlot.DropOnSnapshotRestart = true;
    
    // Auto-create publication for all tables
    o.Publication.Mode = PublicationMode.AllTables;
    
    // Or specify tables
    // o.Publication.Mode = PublicationMode.IncludeTables;
    // o.Publication.IncludedTables.Add("public.users");
    // o.Publication.IncludedTables.Add("public.orders");

    // Or specify schemas
    // o.Publication.Mode = PublicationMode.IncludeSchemas;
    // o.Publication.IncludedSchemas.Add("public");
    // o.Publication.SyncMode = PublicationSyncMode.Exact; // drop extras to match config exactly

    // Source-level connector filters (applied before sink/middleware)
    // o.Filters.Enabled = true;
    // o.Filters.IncludeTables.Add("public.users");
    // o.Filters.ExcludeOperations.Add("u");
});

services.AddSingleton<PostgresCdcProcessor>();
var processor = provider.GetRequiredService<PostgresCdcProcessor>();
await processor.StartAsync(cancellationToken);

Optional multi-sink registration with per-sink filter/middleware:

services.AddSingleton(new SinkRegistration("file", fileSink));
services.AddSingleton(new SinkRegistration("orders-kafka", kafkaSink)
{
    Filters = new SinkFilterOptions
    {
        Enabled = true,
        IncludeTables = ["public.orders"]
    },
    Middlewares =
    [
        new MaskFieldCdcPayloadMiddleware(["credit_card_number"])
    ]
});

Snapshot Strategies

TransactionalExported

  • Exports a consistent snapshot using pg_export_snapshot()
  • Creates replication slot at consistent point
  • Best for: Initial loads where consistency is critical

Incremental

  • Debezium-style chunked snapshot
  • Processes tables in chunks while streaming live WAL changes
  • Merges snapshot and WAL data per-chunk window
  • Best for: Large databases, minimal downtime

BestEffort

  • Simple SELECT snapshot without transaction guarantees
  • Non-blocking, fastest option
  • Best for: Development, non-critical workloads

Configuration

See CdcProcessorOptions in WalFlow.Abstractions for all configuration options.

Replication Slot Lifecycle

  • ReplicationSlot.ManagementMode = EnsureExists (default): create slot if missing, otherwise reuse.
  • ReplicationSlot.ManagementMode = RequireExisting: fail startup when slot does not exist.
  • ReplicationSlot.ManagementMode = RecreateOnStartup: drop and recreate slot once on startup.
  • ReplicationSlot.CleanupMode = Keep (default): keep slot on stop.
  • ReplicationSlot.CleanupMode = DropOnStop: drop slot on graceful stop (similar to Debezium slot.drop.on.stop=true).

Publication Management

  • Publication.Mode = AllTables: create/validate FOR ALL TABLES.
  • Publication.Mode = IncludeTables: manage table members from Publication.IncludedTables.
  • Publication.Mode = IncludeSchemas: manage schema members from Publication.IncludedSchemas.
  • Publication.SyncMode = AddMissing (default): add configured members and keep extras.
  • Publication.SyncMode = Exact: add missing and remove extras to match config exactly.
  • Publication.RecreateOnModeMismatch = true: drop/recreate publication if existing mode differs from configured mode.

Connector Filters

  • Filters.IncludeSchemas / Filters.ExcludeSchemas
  • Filters.IncludeTables / Filters.ExcludeTables (schema.table)
  • Filters.IncludeOperations / Filters.ExcludeOperations (for example c, u, d, r)

These filters are applied in the source before payloads are sent to sinks.

Sink Middleware Layers

WalFlow uses two sink middleware layers:

  • Global pipeline: CdcProcessorOptions.Smt plus injected custom middleware. Runs once before fan-out.
  • Per-sink pipeline: configured on each SinkRegistration. Runs only for that sink.

License

MIT

Product Compatible and additional computed target framework versions.
.NET net10.0 is compatible.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages (6)

Showing the top 5 NuGet packages that depend on WalFlow.Sources.Postgres:

Package Downloads
WalFlow.Failover.Postgres

PostgreSQL lease store and worker adapter plugin for WalFlow failover.

WalFlow.Signals.File

File-based runtime signaling channel plugin for WalFlow incremental snapshot control.

WalFlow.Signals.Kafka

Kafka runtime signaling channel plugin for WalFlow incremental snapshot control.

WalFlow.Signals.AspNetCore

ASP.NET HTTP runtime signaling channel plugin for WalFlow incremental snapshot control.

WalFlow.Signals.Redis

Redis runtime signaling channel plugin for WalFlow incremental snapshot control.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
0.1.0-alpha.alpha... 135 7/8/2026
0.1.0-alpha.alpha... 99 3/30/2026