WalFlow.Sinks.SNS 0.1.0-alpha.alpha.20260708090623

This is a prerelease version of WalFlow.Sinks.SNS.
dotnet add package WalFlow.Sinks.SNS --version 0.1.0-alpha.alpha.20260708090623
                    
NuGet\Install-Package WalFlow.Sinks.SNS -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.Sinks.SNS" 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.Sinks.SNS" Version="0.1.0-alpha.alpha.20260708090623" />
                    
Directory.Packages.props
<PackageReference Include="WalFlow.Sinks.SNS" />
                    
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.Sinks.SNS --version 0.1.0-alpha.alpha.20260708090623
                    
#r "nuget: WalFlow.Sinks.SNS, 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.Sinks.SNS@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.Sinks.SNS&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Addin
#tool nuget:?package=WalFlow.Sinks.SNS&version=0.1.0-alpha.alpha.20260708090623&prerelease
                    
Install as a Cake Tool

WalFlow.Sinks.SNS

AWS SNS (Simple Notification Service) sink for WalFlow CDC processor. Publishes Debezium-format CDC events to SNS topics with message attributes for filtering and fan-out to multiple subscribers.

Features

  • Fan-Out Messaging: Publish once, deliver to multiple subscribers (SQS, Lambda, HTTP, email, SMS)
  • Message Attributes: Automatic attributes for table name and operation type (filtering support)
  • Credential Flexibility: Explicit keys, AWS profiles, or default credential chain (IAM roles)
  • HttpClient Injection: Supports custom HttpClient or IHttpClientFactory for retries/proxies
  • LocalStack Support: Compatible with LocalStack and custom SNS endpoints
  • Optional Message Subject: Configurable subject line for email subscriptions
  • Error Logging: Detailed logging of publish failures

Installation

dotnet add package WalFlow.Sinks.SNS

Configuration Options

Option Type Default Description
TopicArn string required SNS topic ARN
Region string "us-east-1" AWS region
ServiceUrl string? null Custom endpoint (for LocalStack)
AccessKeyId string? null AWS access key (optional)
SecretAccessKey string? null AWS secret key (optional)
ProfileName string? null AWS profile name (optional)
MessageSubject string? null Subject for email notifications

Usage

Basic Registration with Explicit Credentials

using Microsoft.Extensions.DependencyInjection;
using WalFlow.Abstractions;
using WalFlow.Sinks.SNS;

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.AccessKeyId = "AKIAIOSFODNN7EXAMPLE";
    opt.SecretAccessKey = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
    opt.MessageSubject = "WalFlow CDC Event";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    new HttpClient(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));

Using AWS Profile

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-west-2:123456789012:walflow-cdc-events";
    opt.Region = "us-west-2";
    opt.ProfileName = "my-aws-profile";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    new HttpClient(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));

Using Default Credential Chain (IAM Roles)

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events";
    opt.Region = "us-east-1";
    // No credentials specified - will use IAM role, environment variables, etc.
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    new HttpClient(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));
services.AddHttpClient("SNS")
    .ConfigureHttpClient(client =>
    {
        client.Timeout = TimeSpan.FromSeconds(30);
    })
    .AddPolicyHandler(GetRetryPolicy());  // Polly retry policy

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.MessageSubject = "CDC Event";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClientFactory(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    sp.GetRequiredService<IHttpClientFactory>(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));

static IAsyncPolicy<HttpResponseMessage> GetRetryPolicy()
{
    return HttpPolicyExtensions
        .HandleTransientHttpError()
        .WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromSeconds(Math.Pow(2, retryAttempt)));
}

Using Custom AmazonSimpleNotificationServiceClient

var snsClient = new AmazonSimpleNotificationServiceClient(
    new BasicAWSCredentials("access-key", "secret-key"),
    new AmazonSimpleNotificationServiceConfig
    {
        RegionEndpoint = RegionEndpoint.USEast1,
        MaxErrorRetry = 3,
        Timeout = TimeSpan.FromSeconds(30)
    }
);

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events";
    opt.MessageSubject = "CDC Event";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithSNSClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    snsClient,
    sp.GetRequiredService<ILogger<SNSSink>>()
));

LocalStack Configuration

services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:000000000000:walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.ServiceUrl = "http://localhost:4566";  // LocalStack
    opt.AccessKeyId = "test";
    opt.SecretAccessKey = "test";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    new HttpClient(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));

Direct Instantiation

var options = Options.Create(new SNSSinkOptions
{
    TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events",
    Region = "us-east-1",
    MessageSubject = "CDC Event"
});

var logger = loggerFactory.CreateLogger<SNSSink>();
var sink = SNSSink.CreateWithHttpClient(options, new HttpClient(), logger);

// Use with CDC processor
var processor = new PostgresCdcProcessor(cdcOptions, sink, stateStore, processorLogger);
await processor.StartAsync(cancellationToken);

// Cleanup
sink.Dispose();

Message Attributes

The sink automatically adds the following message attributes to each published event:

Attribute Type Description Example
table String Fully qualified table name "public.users"
operation String CDC operation type "INSERT", "UPDATE", "DELETE"

These attributes enable subscription filtering without needing to parse the message body.

Subscription Filter Example

Filter to only receive INSERT operations:

{
  "operation": ["INSERT"]
}

Filter to only receive events from the users table:

{
  "table": ["public.users"]
}

Filter for UPDATE or DELETE on the orders table:

{
  "table": ["public.orders"],
  "operation": ["UPDATE", "DELETE"]
}

Creating an SNS Topic

Using AWS CLI

aws sns create-topic --name walflow-cdc-events

Subscribe SQS Queue

aws sns subscribe \
  --topic-arn arn:aws:sns:us-east-1:123456789012:walflow-cdc-events \
  --protocol sqs \
  --notification-endpoint arn:aws:sqs:us-east-1:123456789012:my-queue

Subscribe Lambda Function

aws sns subscribe \
  --topic-arn arn:aws:sns:us-east-1:123456789012:walflow-cdc-events \
  --protocol lambda \
  --notification-endpoint arn:aws:lambda:us-east-1:123456789012:function:process-cdc

Subscribe HTTP Endpoint

aws sns subscribe \
  --topic-arn arn:aws:sns:us-east-1:123456789012:walflow-cdc-events \
  --protocol https \
  --notification-endpoint https://api.example.com/webhook/cdc

Message Format

The message body contains the full Debezium-format CDC event as JSON:

{
  "source": {
    "connector": "walflow-postgres",
    "name": "walflow_slot",
    "ts_ms": 1672531200000,
    "snapshot": false,
    "db": "mydb",
    "schema": "public",
    "table": "users",
    "lsn": 12345
  },
  "op": "c",
  "ts_ms": 1672531200000,
  "before": null,
  "after": {
    "id": 1,
    "name": "Alice",
    "email": "alice@example.com"
  }
}

Publishing Behavior

  • Synchronous: Events are published immediately (no batching)
  • Non-blocking: Publish failures are logged but don't throw exceptions
  • At-Least-Once: SNS guarantees at-least-once delivery to all subscribers
  • No Ordering: SNS does not guarantee message order

Why No Batching?

Unlike Kinesis and SQS, SNS does not support batch publishing. Each event is published individually via PublishAsync().

Error Handling

Failed publishes are logged with details but do NOT throw exceptions. Check logs for:

Failed to publish message to SNS: {ErrorMessage}
Topic: {topicArn}
Table: {table}
Operation: {operation}
Payload: {serializedPayload}

Common error codes:

  • NotFound: Topic ARN is invalid or topic was deleted
  • InvalidParameter: Malformed message or attributes
  • AuthorizationError: Insufficient permissions

AWS Credentials Priority

The sink uses the following credential priority:

  1. Explicit credentials (AccessKeyId + SecretAccessKey)
  2. AWS Profile (ProfileName)
  3. Default credential chain:
    • Environment variables (AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY)
    • IAM role (EC2, ECS, Lambda)
    • ~/.aws/credentials file

Fan-Out Architecture

SNS is ideal for broadcasting CDC events to multiple downstream systems:

PostgresCdcProcessor
        ↓
    SNS Topic
        ↓
    ┌───────┼───────┐
    ↓       ↓       ↓
  SQS   Lambda   HTTP
   ↓       ↓       ↓
Analytics Search  API

Example: Multi-Subscriber Setup

// Configure SNS sink
services.Configure<SNSSinkOptions>(opt =>
{
    opt.TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events";
    opt.Region = "us-east-1";
});

services.AddSingleton<ISink>(sp => SNSSink.CreateWithHttpClient(
    sp.GetRequiredService<IOptions<SNSSinkOptions>>(),
    new HttpClient(),
    sp.GetRequiredService<ILogger<SNSSink>>()
));

Then create subscriptions with filters:

  1. Analytics Queue (all operations):

    aws sns subscribe --topic-arn ... --protocol sqs --notification-endpoint arn:aws:sqs:...:analytics-queue
    
  2. Search Index Lambda (INSERT/UPDATE only):

    aws sns subscribe --topic-arn ... --protocol lambda --notification-endpoint arn:aws:lambda:...:index-docs \
      --attributes '{"FilterPolicy":"{\"operation\":[\"INSERT\",\"UPDATE\"]}"}'
    
  3. Audit Webhook (DELETE only):

    aws sns subscribe --topic-arn ... --protocol https --notification-endpoint https://audit.example.com/webhook \
      --attributes '{"FilterPolicy":"{\"operation\":[\"DELETE\"]}"}'
    

Monitoring

Key metrics to monitor:

  • Publish Success Rate: Check CloudWatch NumberOfMessagesPublished metric
  • Notification Delivery: Monitor NumberOfNotificationsFailed per subscription
  • Fan-Out Lag: Monitor downstream consumer metrics (SQS queue depth, Lambda errors)

Comparison with Other Sinks

Feature SNS Kinesis SQS File
Batching ❌ ✅ 500 ✅ 10 ❌
Ordering ❌ ✅ Per-shard ✅ FIFO only ✅
Throughput ⚡ High ⚡ High 🔸 Medium ⚡ High
Replay ❌ ✅ 24h-365d ❌ ✅
Fanout ✅ ✅ Multi-consumer ❌ ❌
Filtering ✅ Attributes ❌ ❌ ❌
Use Case Broadcast Streaming Queueing Debug

Best Practices

  1. Use Message Attributes: Leverage table and operation attributes for subscription filtering
  2. Monitor Failed Deliveries: Set up CloudWatch alarms for NumberOfNotificationsFailed
  3. Dead-Letter Queues: Configure DLQs on SQS subscriptions for failed deliveries
  4. Idempotent Consumers: SNS provides at-least-once delivery, so consumers must handle duplicates
  5. Topic Policies: Use IAM policies to control who can publish/subscribe

License

See the LICENSE file in the repository root.

Inline DI Registration

services.AddSingleton<ISink>(sp =>
{
    var options = Options.Create(new SNSSinkOptions
    {
        TopicArn = "arn:aws:sns:us-east-1:123456789012:walflow-cdc-events",
        Region = "us-east-1"
    });

    return SNSSink.CreateWithHttpClient(
        options,
        sp.GetRequiredService<ILogger<SNSSink>>(),
        sp.GetRequiredService<HttpClient>());
});
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 (1)

Showing the top 1 NuGet packages that depend on WalFlow.Sinks.SNS:

Package Downloads
WalFlow

Meta-package that pulls in the full WalFlow CDC stack for consumers who want the default everything-included setup.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
0.1.0-alpha.alpha... 85 7/8/2026
0.1.0-alpha.alpha... 81 3/30/2026