WalFlow.Sinks.Kinesis
0.1.0-alpha.alpha.20260708090623
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
<PackageReference Include="WalFlow.Sinks.Kinesis" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageVersion Include="WalFlow.Sinks.Kinesis" Version="0.1.0-alpha.alpha.20260708090623" />
<PackageReference Include="WalFlow.Sinks.Kinesis" />
paket add WalFlow.Sinks.Kinesis --version 0.1.0-alpha.alpha.20260708090623
#r "nuget: WalFlow.Sinks.Kinesis, 0.1.0-alpha.alpha.20260708090623"
#:package WalFlow.Sinks.Kinesis@0.1.0-alpha.alpha.20260708090623
#addin nuget:?package=WalFlow.Sinks.Kinesis&version=0.1.0-alpha.alpha.20260708090623&prerelease
#tool nuget:?package=WalFlow.Sinks.Kinesis&version=0.1.0-alpha.alpha.20260708090623&prerelease
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
HttpClientorIHttpClientFactoryfor 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>>()
));
Using IHttpClientFactory (Recommended)
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:
- Events are buffered in a
BlockingCollection - When either
MaxBatchSize(500) orBatchingWindowMs(1000ms) is reached, a batch is flushed - 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 rateInvalidArgumentException: Malformed data or partition keyResourceNotFoundException: Stream doesn't exist
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 (
Monitoring
Key metrics to monitor:
- PutRecords Success Rate: Check CloudWatch
PutRecords.Successmetric - Iterator Age: High age indicates consumer lag
- Write Throughput: Monitor
IncomingBytesandIncomingRecords - 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 | 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.Kinesis (>= 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.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 |