WalFlow.Sources.Postgres
0.1.0-alpha.alpha.20260708090623
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
<PackageReference Include="WalFlow.Sources.Postgres" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageVersion Include="WalFlow.Sources.Postgres" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageReference Include="WalFlow.Sources.Postgres" />
paket add WalFlow.Sources.Postgres --version 0.1.0-alpha.alpha.20260708090623
#r "nuget: WalFlow.Sources.Postgres, 0.1.0-alpha.alpha.20260708090623"
#:package WalFlow.Sources.Postgres@0.1.0-alpha.alpha.20260708090623
#addin nuget:?package=WalFlow.Sources.Postgres&version=0.1.0-alpha.alpha.20260708090623&prerelease
#tool nuget:?package=WalFlow.Sources.Postgres&version=0.1.0-alpha.alpha.20260708090623&prerelease
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
SELECTsnapshot 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 Debeziumslot.drop.on.stop=true).
Publication Management
Publication.Mode = AllTables: create/validateFOR ALL TABLES.Publication.Mode = IncludeTables: manage table members fromPublication.IncludedTables.Publication.Mode = IncludeSchemas: manage schema members fromPublication.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.ExcludeSchemasFilters.IncludeTables/Filters.ExcludeTables(schema.table)Filters.IncludeOperations/Filters.ExcludeOperations(for examplec,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.Smtplus injected custom middleware. Runs once before fan-out. - Per-sink pipeline: configured on each
SinkRegistration. Runs only for that sink.
License
MIT
| Product | Versions 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. |
-
net10.0
- Microsoft.Extensions.DependencyInjection (>= 10.0.9)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Logging (>= 10.0.9)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Options (>= 10.0.9)
- Npgsql (>= 10.0.3)
- WalFlow.Abstractions (>= 0.1.0-alpha.alpha.20260708090623)
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 |