WalFlow.Sinks.Kinesis 0.1.0-alpha.alpha.20260708090623

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

WalFlow.Sinks.Kinesis

AWS Kinesis Data Streams sink for WalFlow CDC processor. Publishes Debezium-format CDC events to Kinesis with automatic batching and configurable partition key strategies.

Features

  • Automatic Batching: Buffers records (max 500, 1000ms window) for efficient PutRecords calls
  • Partition Key Strategies: SlotName, table-based, random, or fixed partition keys
  • 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 Kinesis endpoints
  • Graceful Shutdown: Flushes buffered records on disposal
  • Error Logging: Detailed logging of PutRecords failures per record

Installation

dotnet add package WalFlow.Sinks.Kinesis

Configuration Options

Option Type Default Description
StreamName string required Kinesis stream name
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)
PartitionKeyStrategy PartitionKeyStrategy SlotName Partition key strategy
FixedPartitionKey string? null Used when strategy is Fixed
MaxBatchSize int 500 Max records per PutRecords call
BatchingWindowMs int 1000 Max time to buffer before flush

Partition Key Strategies

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
Random Uses random GUID Maximum throughput, no ordering
Fixed Uses FixedPartitionKey option Custom partition logic

Usage

Basic Registration with Explicit Credentials

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

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.AccessKeyId = "AKIAIOSFODNN7EXAMPLE";
    opt.SecretAccessKey = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
    opt.PartitionKeyStrategy = PartitionKeyStrategy.Table;
});

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

Using AWS Profile

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.Region = "us-west-2";
    opt.ProfileName = "my-aws-profile";
    opt.PartitionKeyStrategy = PartitionKeyStrategy.SlotName;
});

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

Using Default Credential Chain (IAM Roles)

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.Region = "us-east-1";
    // No credentials specified - will use IAM role, environment variables, etc.
    opt.PartitionKeyStrategy = PartitionKeyStrategy.Table;
});

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

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.PartitionKeyStrategy = PartitionKeyStrategy.Table;
});

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

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

Using Custom AmazonKinesisClient

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

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.PartitionKeyStrategy = PartitionKeyStrategy.Table;
});

services.AddSingleton<ISink>(sp => KinesisSink.CreateWithKinesisClient(
    sp.GetRequiredService<IOptions<KinesisSinkOptions>>(),
    kinesisClient,
    sp.GetRequiredService<ILogger<KinesisSink>>()
));

LocalStack Configuration

services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.Region = "us-east-1";
    opt.ServiceUrl = "http://localhost:4566";  // LocalStack
    opt.AccessKeyId = "test";
    opt.SecretAccessKey = "test";
    opt.PartitionKeyStrategy = PartitionKeyStrategy.SlotName;
});

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

Direct Instantiation

var options = Options.Create(new KinesisSinkOptions
{
    StreamName = "walflow-cdc-events",
    Region = "us-east-1",
    PartitionKeyStrategy = PartitionKeyStrategy.Table,
    MaxBatchSize = 500,
    BatchingWindowMs = 1000
});

var logger = loggerFactory.CreateLogger<KinesisSink>();
var sink = KinesisSink.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 records
sink.Dispose();

Partition Key Strategies Explained

SlotName (Default)

All events from the same replication slot go to the same partition. Preserves total order of all changes.

opt.PartitionKeyStrategy = PartitionKeyStrategy.SlotName;

Result: partition_key = "walflow_slot"

Table

Each table gets its own partition key. Preserves per-table order, allows parallel processing of different tables.

opt.PartitionKeyStrategy = PartitionKeyStrategy.Table;

Result: partition_key = "public.users" for users table, "public.orders" for orders table

Random

Each event gets a random partition key. Maximum throughput, no ordering guarantees.

opt.PartitionKeyStrategy = PartitionKeyStrategy.Random;

Result: partition_key = "a3f2b4c1-..." (random GUID)

Fixed

All events use a custom fixed partition key.

opt.PartitionKeyStrategy = PartitionKeyStrategy.Fixed;
opt.FixedPartitionKey = "my-custom-key";

Result: partition_key = "my-custom-key"

Batching Behavior

The sink uses background batching to optimize PutRecords calls:

  1. Events are buffered in a BlockingCollection
  2. When either MaxBatchSize (500) or BatchingWindowMs (1000ms) is reached, a batch is flushed
  3. On disposal, all buffered events are flushed before shutdown
services.Configure<KinesisSinkOptions>(opt =>
{
    opt.StreamName = "walflow-cdc-events";
    opt.MaxBatchSize = 250;        // Smaller batches
    opt.BatchingWindowMs = 500;     // Flush more frequently
});

Error Handling

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

Failed to write record: {ErrorCode} - {ErrorMessage}
PartitionKey: {partitionKey}
Payload: {serializedPayload}

Common error codes:

  • ProvisionedThroughputExceededException: Increase shard count or reduce write rate
  • InvalidArgumentException: Malformed data or partition key
  • ResourceNotFoundException: Stream doesn't exist

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

Monitoring

Key metrics to monitor:

  • PutRecords Success Rate: Check CloudWatch PutRecords.Success metric
  • Iterator Age: High age indicates consumer lag
  • Write Throughput: Monitor IncomingBytes and IncomingRecords
  • Batch Size: Larger batches = better efficiency

Comparison with Other Sinks

Feature Kinesis SQS SNS File
Batching ✅ 500 ✅ 10 ❌ ❌
Ordering ✅ Per-shard ✅ FIFO only ❌ ✅
Throughput ⚡ High 🔸 Medium ⚡ High ⚡ High
Replay ✅ 24h-365d ❌ ❌ ✅
Fanout ✅ Multi-consumer ❌ ✅ ❌

License

See the LICENSE file in the repository root.

Inline DI Registration

services.AddSingleton<ISink>(sp =>
{
    var options = Options.Create(new KinesisSinkOptions
    {
        StreamName = "walflow-cdc-events",
        Region = "us-east-1"
    });

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

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... 85 3/30/2026