WalFlow.Sinks.SNS
0.1.0-alpha.alpha.20260708090623
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
<PackageReference Include="WalFlow.Sinks.SNS" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageVersion Include="WalFlow.Sinks.SNS" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageReference Include="WalFlow.Sinks.SNS" />
paket add WalFlow.Sinks.SNS --version 0.1.0-alpha.alpha.20260708090623
#r "nuget: WalFlow.Sinks.SNS, 0.1.0-alpha.alpha.20260708090623"
#:package WalFlow.Sinks.SNS@0.1.0-alpha.alpha.20260708090623
#addin nuget:?package=WalFlow.Sinks.SNS&version=0.1.0-alpha.alpha.20260708090623&prerelease
#tool nuget:?package=WalFlow.Sinks.SNS&version=0.1.0-alpha.alpha.20260708090623&prerelease
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
HttpClientorIHttpClientFactoryfor 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>>()
));
Using IHttpClientFactory (Recommended)
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 deletedInvalidParameter: Malformed message or attributesAuthorizationError: Insufficient permissions
AWS Credentials Priority
The sink uses the following credential priority:
- Explicit credentials (
AccessKeyId+SecretAccessKey) - AWS Profile (
ProfileName) - Default credential chain:
- Environment variables (
AWS_ACCESS_KEY_ID,AWS_SECRET_ACCESS_KEY) - IAM role (EC2, ECS, Lambda)
~/.aws/credentialsfile
- Environment variables (
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:
Analytics Queue (all operations):
aws sns subscribe --topic-arn ... --protocol sqs --notification-endpoint arn:aws:sqs:...:analytics-queueSearch Index Lambda (INSERT/UPDATE only):
aws sns subscribe --topic-arn ... --protocol lambda --notification-endpoint arn:aws:lambda:...:index-docs \ --attributes '{"FilterPolicy":"{\"operation\":[\"INSERT\",\"UPDATE\"]}"}'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
NumberOfMessagesPublishedmetric - Notification Delivery: Monitor
NumberOfNotificationsFailedper 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
- Use Message Attributes: Leverage
tableandoperationattributes for subscription filtering - Monitor Failed Deliveries: Set up CloudWatch alarms for
NumberOfNotificationsFailed - Dead-Letter Queues: Configure DLQs on SQS subscriptions for failed deliveries
- Idempotent Consumers: SNS provides at-least-once delivery, so consumers must handle duplicates
- 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 | 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
- AWSSDK.Core (>= 4.0.100.1)
- AWSSDK.SimpleNotificationService (>= 4.0.100.1)
- Microsoft.Extensions.Configuration (>= 10.0.9)
- Microsoft.Extensions.Configuration.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Configuration.Binder (>= 10.0.9)
- Microsoft.Extensions.DependencyInjection (>= 10.0.9)
- Microsoft.Extensions.DependencyInjection.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Diagnostics.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Http (>= 10.0.9)
- Microsoft.Extensions.Logging (>= 10.0.9)
- Microsoft.Extensions.Logging.Abstractions (>= 10.0.9)
- Microsoft.Extensions.Options (>= 10.0.9)
- Microsoft.Extensions.Options.ConfigurationExtensions (>= 10.0.9)
- WalFlow.Abstractions (>= 0.1.0-alpha.alpha.20260708090623)
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 |