WalFlow.Sinks.SQS
0.1.0-alpha.alpha.20260708090623
dotnet add package WalFlow.Sinks.SQS --version 0.1.0-alpha.alpha.20260708090623
NuGet\Install-Package WalFlow.Sinks.SQS -Version 0.1.0-alpha.alpha.20260708090623
<PackageReference Include="WalFlow.Sinks.SQS" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageVersion Include="WalFlow.Sinks.SQS" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageReference Include="WalFlow.Sinks.SQS" />
paket add WalFlow.Sinks.SQS --version 0.1.0-alpha.alpha.20260708090623
#r "nuget: WalFlow.Sinks.SQS, 0.1.0-alpha.alpha.20260708090623"
#:package WalFlow.Sinks.SQS@0.1.0-alpha.alpha.20260708090623
#addin nuget:?package=WalFlow.Sinks.SQS&version=0.1.0-alpha.alpha.20260708090623&prerelease
#tool nuget:?package=WalFlow.Sinks.SQS&version=0.1.0-alpha.alpha.20260708090623&prerelease
WalFlow.Sinks.SQS
AWS SQS (Simple Queue Service) sink for WalFlow CDC processor. Publishes Debezium-format CDC events to SQS with automatic batching, FIFO queue support, and configurable message grouping.
Features
- Automatic Batching: Buffers messages (max 10, 1000ms window) for efficient SendMessageBatch calls
- FIFO Queue Support: Preserves ordering with configurable MessageGroupId strategies
- Standard Queue Support: High throughput without ordering guarantees
- 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 SQS endpoints
- Graceful Shutdown: Flushes buffered messages on disposal
- Error Logging: Detailed logging of SendMessageBatch failures per message
Installation
dotnet add package WalFlow.Sinks.SQS
Configuration Options
| Option | Type | Default | Description |
|---|---|---|---|
QueueUrl |
string |
required | SQS queue URL |
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) |
MessageGroupIdStrategy |
MessageGroupIdStrategy |
SlotName |
Message group ID strategy (FIFO only) |
FixedMessageGroupId |
string? |
null |
Used when strategy is Fixed (FIFO only) |
MaxBatchSize |
int |
10 |
Max messages per SendMessageBatch call |
BatchingWindowMs |
int |
1000 |
Max time to buffer before flush |
UseMessageDeduplicationId |
bool |
false |
Enables FIFO MessageDeduplicationId values |
Message Group ID Strategies (FIFO Queues Only)
| Strategy | Description | Use Case |
|---|---|---|
SlotName |
Uses replication slot name | Single-table or grouped changes |
Table |
Uses {schema}.{table} |
Distribute by table, preserve per-table order |
Fixed |
Uses FixedMessageGroupId option |
Custom grouping logic |
Note: MessageGroupId is only used for FIFO queues (URLs ending in .fifo). For standard queues, it's ignored.
FIFO Deduplication IDs
When UseMessageDeduplicationId is enabled, SQS sends a FIFO MessageDeduplicationId.
By default WalFlow derives the ID from a SHA-256 hash of the transformed message body.
Pass a code-defined ISinkIdempotencyKeyResolver to the sink factory when the deduplication key should come from a stable event identity instead of message formatting:
using Amazon.SQS;
using WalFlow.Abstractions.Sinks;
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events.fifo";
opt.MessageGroupId = "walflow-cdc";
opt.UseMessageDeduplicationId = true;
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithSQSClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
sp.GetRequiredService<ILogger<SQSSink>>(),
sp.GetRequiredService<IAmazonSQS>(),
messageDeduplicationIdResolver: SinkIdempotencyKeyResolvers.FromDelegate(payload =>
$"{payload.Source?.Name}:{payload.Source?.Schema}.{payload.Source?.Table}:{payload.Source?.Lsn}:{payload.Operation}:{payload.Transaction?.TotalOrder}:{payload.After?["id"] ?? payload.Before?["id"]}")));
Resolved IDs must identify one logical event within SQS's deduplication window, use SQS's allowed ASCII character set, and be no longer than the 128-character provider limit.
Usage
Basic Registration with Standard Queue
using Microsoft.Extensions.DependencyInjection;
using WalFlow.Abstractions;
using WalFlow.Sinks.SQS;
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events";
opt.Region = "us-east-1";
opt.AccessKeyId = "AKIAIOSFODNN7EXAMPLE";
opt.SecretAccessKey = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
new HttpClient(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
FIFO Queue with Table-Based Grouping
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events.fifo";
opt.Region = "us-east-1";
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.Table; // One group per table
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
new HttpClient(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
Using AWS Profile
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-west-2.amazonaws.com/123456789012/walflow-cdc-events";
opt.Region = "us-west-2";
opt.ProfileName = "my-aws-profile";
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
new HttpClient(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
Using Default Credential Chain (IAM Roles)
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events";
opt.Region = "us-east-1";
// No credentials specified - will use IAM role, environment variables, etc.
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
new HttpClient(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
Using IHttpClientFactory (Recommended)
services.AddHttpClient("SQS")
.ConfigureHttpClient(client =>
{
client.Timeout = TimeSpan.FromSeconds(30);
})
.AddPolicyHandler(GetRetryPolicy()); // Polly retry policy
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events.fifo";
opt.Region = "us-east-1";
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.SlotName;
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClientFactory(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
sp.GetRequiredService<IHttpClientFactory>(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
static IAsyncPolicy<HttpResponseMessage> GetRetryPolicy()
{
return HttpPolicyExtensions
.HandleTransientHttpError()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromSeconds(Math.Pow(2, retryAttempt)));
}
Using Custom AmazonSQSClient
var sqsClient = new AmazonSQSClient(
new BasicAWSCredentials("access-key", "secret-key"),
new AmazonSQSConfig
{
RegionEndpoint = RegionEndpoint.USEast1,
MaxErrorRetry = 3,
Timeout = TimeSpan.FromSeconds(30)
}
);
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events.fifo";
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.Table;
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithSQSClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
sqsClient,
sp.GetRequiredService<ILogger<SQSSink>>()
));
LocalStack Configuration
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "http://localhost:4566/000000000000/walflow-cdc-events.fifo";
opt.Region = "us-east-1";
opt.ServiceUrl = "http://localhost:4566"; // LocalStack
opt.AccessKeyId = "test";
opt.SecretAccessKey = "test";
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.SlotName;
});
services.AddSingleton<ISink>(sp => SQSSink.CreateWithHttpClient(
sp.GetRequiredService<IOptions<SQSSinkOptions>>(),
new HttpClient(),
sp.GetRequiredService<ILogger<SQSSink>>()
));
Direct Instantiation
var options = Options.Create(new SQSSinkOptions
{
QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events.fifo",
Region = "us-east-1",
MessageGroupIdStrategy = MessageGroupIdStrategy.Table,
MaxBatchSize = 10,
BatchingWindowMs = 1000
});
var logger = loggerFactory.CreateLogger<SQSSink>();
var sink = SQSSink.CreateWithHttpClient(options, new HttpClient(), logger);
// Use with CDC processor
var processor = new PostgresCdcProcessor(cdcOptions, sink, stateStore, processorLogger);
await processor.StartAsync(cancellationToken);
// Graceful shutdown - flushes buffered messages
sink.Dispose();
FIFO Queue Setup
Create FIFO Queue
aws sqs create-queue \
--queue-name walflow-cdc-events.fifo \
--attributes FifoQueue=true,ContentBasedDeduplication=true
Important: FIFO queue names MUST end with .fifo.
Message Deduplication
When UseMessageDeduplicationId is enabled, the sink generates MessageDeduplicationId using SHA-256 of the message body unless a code-defined resolver is supplied. This prevents duplicate messages with the same deduplication key within the SQS 5-minute deduplication window.
Message Group ID Strategies
SlotName (Default)
All events from the same replication slot go to the same message group. Preserves total order.
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.SlotName;
Result: MessageGroupId = "walflow_slot"
Table
Each table gets its own message group. Preserves per-table order, allows parallel processing.
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.Table;
Result:
MessageGroupId = "public.users"for users tableMessageGroupId = "public.orders"for orders table
Fixed
All events use a custom fixed message group ID.
opt.MessageGroupIdStrategy = MessageGroupIdStrategy.Fixed;
opt.FixedMessageGroupId = "my-custom-group";
Result: MessageGroupId = "my-custom-group"
Batching Behavior
The sink uses background batching to optimize SendMessageBatch calls:
- Messages are buffered in a
BlockingCollection - When either
MaxBatchSize(10) orBatchingWindowMs(1000ms) is reached, a batch is flushed - On disposal, all buffered messages are flushed before shutdown
services.Configure<SQSSinkOptions>(opt =>
{
opt.QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events";
opt.MaxBatchSize = 5; // Smaller batches
opt.BatchingWindowMs = 500; // Flush more frequently
});
Note: SQS batch limit is 10 messages per request (AWS limitation).
Error Handling
Failed messages are logged with details but do NOT throw exceptions. Check logs for:
Failed to send message: {ErrorCode} - {ErrorMessage}
MessageId: {messageId}
Payload: {serializedPayload}
Common error codes:
QueueDoesNotExist: Queue URL is invalidInvalidMessageContents: Message body is malformedBatchRequestTooLong: Batch exceeds 256KB (reduce batch size)AWS.SimpleQueueService.NonExistentQueue: Queue was deleted
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 (
Standard vs FIFO Queues
| Feature | Standard | FIFO |
|---|---|---|
| Throughput | Unlimited | 300 TPS (or 3000 with batching) |
| Ordering | Best-effort | Strict per MessageGroupId |
| Deduplication | No | Yes (5 min window) |
| Message Group | N/A | Required |
| Use Case | High throughput, at-least-once | Ordered processing, exactly-once |
Monitoring
Key metrics to monitor:
- SendMessageBatch Success Rate: Check CloudWatch
NumberOfMessagesSentmetric - Queue Depth: Monitor
ApproximateNumberOfMessagesVisible - Age of Oldest Message: High age indicates consumer lag
- Batch Size: Larger batches = better efficiency (up to 10 for SQS)
Comparison with Other Sinks
| Feature | SQS | Kinesis | SNS | File |
|---|---|---|---|---|
| Batching | ✅ 10 | ✅ 500 | ❌ | ❌ |
| Ordering | ✅ FIFO only | ✅ Per-shard | ❌ | ✅ |
| Throughput | 🔸 Medium | ⚡ High | ⚡ High | ⚡ High |
| Replay | ❌ | ✅ 24h-365d | ❌ | ✅ |
| Fanout | ❌ | ✅ Multi-consumer | ✅ | ❌ |
| Max Retention | 14 days | 365 days | N/A | Unlimited |
License
See the LICENSE file in the repository root.
Inline DI Registration
services.AddSingleton<ISink>(sp =>
{
var options = Options.Create(new SQSSinkOptions
{
QueueUrl = "https://sqs.us-east-1.amazonaws.com/123456789012/walflow-cdc-events",
Region = "us-east-1"
});
return SQSSink.CreateWithHttpClient(
options,
sp.GetRequiredService<ILogger<SQSSink>>(),
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.SQS (>= 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.SQS:
| 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... | 89 | 7/8/2026 |
| 0.1.0-alpha.alpha... | 82 | 3/30/2026 |