WalFlow.Sinks.SQS 0.1.0-alpha.alpha.20260708090623

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

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 HttpClient or IHttpClientFactory for 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>>()
));
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 table
  • MessageGroupId = "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:

  1. Messages are buffered in a BlockingCollection
  2. When either MaxBatchSize (10) or BatchingWindowMs (1000ms) is reached, a batch is flushed
  3. 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 invalid
  • InvalidMessageContents: Message body is malformed
  • BatchRequestTooLong: Batch exceeds 256KB (reduce batch size)
  • AWS.SimpleQueueService.NonExistentQueue: Queue was deleted

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

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 NumberOfMessagesSent metric
  • 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 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.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