Custom Sinks
A sink is a destination plugin. Implement ISink to deliver batches of records anywhere - another database, a message broker, a cache.
TIP
We're always looking for new sink contributions. Feel free to open a pull request for review.
Interface
public interface ISink
{
string Name { get; }
Task<DeliveryResult> DeliverAsync(SinkBatch batch, CancellationToken ct);
}
public sealed record SinkBatch(string SinkName, IReadOnlyList<SinkRecord> Records);
public sealed record SinkRecord(
string? Destination, // e.g. index/topic/table; null = sink default
string DocumentId, // stable id for upsert/delete
IReadOnlyDictionary<string, object?>? Document, // the field bag; null when IsDeletion
bool IsDeletion,
ChangeMetadata Metadata); // source provenanceRecords arrive in commit order. Each is either an upsert of Document under DocumentId, or a deletion of DocumentId.
SinkBatch, SinkRecord, SinkPurgeRequest, ChangeMetadata, ChangeEvent and BackfillState are positional records you may construct yourself (tests, adapters). Their positional parameters are fixed for 1.x: anything Wallaby adds later arrives as an init property with a default, so existing constructor calls keep compiling and a sink that ignores a new property keeps working.
Returning a result
Classify the outcome so the dispatcher can react:
return DeliveryResult.Success; // batch accepted
return DeliveryResult.Retry("503 from upstream"); // transient - retried with backoff
return DeliveryResult.Permanent("schema rejected"); // non-retryable - halts the pipelineRetryable failures are retried with exponential backoff and jitter. A permanent failure (or exhausted retries) halts the pipeline; the batch is retried after the leader session restarts (with its own backoff), so a batch is never silently dropped. The halt is pipeline-wide: every sink shares one replication slot and one acknowledgement point, so no sink receives further batches until the failing one accepts its batch. To isolate a destination whose reliability differs from the others, run it in its own Wallaby worker with its own SlotName and PublicationName. The sink outage walkthrough shows the halt and recovery step by step.
Throw WallabyConfigurationException for a configuration error (for example a record with no resolvable destination); any other exception thrown from DeliverAsync is treated as a permanent failure. Cancellation of ct propagates as-is, and a retryable result returned while ct is cancelled is treated as cancellation, so a classifying catch-all needs no cancellation guard.
Idempotency & ordering
Delivery is at-least-once: the replication slot only advances after a batch is durably delivered, so a crash can redeliver the last batch. Your sink should make delivery idempotent by supporting upsert and delete by DocumentId.
Sinks should also preserve commit order - if you create batches internally, ensure you preserve it.
One-time setup
If your sink needs setup before first delivery (create a topic, configure an index), implement ISinkInitializer:
public sealed class MySink : ISink, ISinkInitializer
{
public Task InitializeAsync(CancellationToken ct) => /* idempotent setup */;
}InitializeAsync runs on the leader, once, after self-config and before streaming begins and again whenever a standby takes over leadership. Make it idempotent. If it throws, the leader session is retried (the pipeline won't stream into an unconfigured sink).
Purging
To support purge-then-backfill (emptying a destination so a fresh backfill converges it to exactly the current table contents), implement ISinkPurger:
public sealed class MySink : ISink, ISinkPurger
{
public Task PurgeAsync(SinkPurgeRequest request, CancellationToken ct)
=> /* delete every document at request.Destination (null = the sink's default destination) */;
}PurgeAsync runs on the leader right before the backfill's snapshot read. Make it idempotent; throw to fail the backfill run (the purge re-runs on the retry). A sink without the capability is skipped with a warning when a purge is requested.
Cleanup
Registering a sink makes Wallaby responsible for its lifetime. If your sink holds resources (a client, a connection pool, a producer), implement IAsyncDisposable (or IDisposable):
public sealed class MySink : ISink, IAsyncDisposable
{
public ValueTask DisposeAsync() => _client.DisposeAsync();
}Disposal runs once, at host shutdown, after streaming has stopped. A sink implementing both interfaces is disposed via DisposeAsync only, and a throwing dispose is logged without disrupting the rest of shutdown.
Registering
// An instance:
cdc.AddSink(new MySink(...))
.WithMappings(sink => sink.Map<Product>().UsingTransform(/* ... */));
// Resolved from the container:
cdc.AddSink("my-sink", sp => new MySink(sp.GetRequiredService<HttpClient>()))
.WithMappings(sink => sink.Map<Product>().UsingTransform(/* ... */));AddSink returns a sink-scoped builder: declare the entities the sink receives in WithMappings(...), or continue the chain via its Wallaby property for a sink registered without mappings.
Envelope helpers
If your sink emits a JSON envelope, SinkEnvelopeJson provides the record-level pieces the built-in HTTP and Kafka sinks share: reflection-free (AOT-safe) document/metadata writing and a stable deduplication key:
using var writer = new Utf8JsonWriter(buffer);
writer.WriteStartObject();
writer.WriteString("id", record.DocumentId);
writer.WriteString("idempotencyKey", SinkEnvelopeJson.IdempotencyKey(record));
writer.WritePropertyName("document");
SinkEnvelopeJson.WriteDocument(writer, record.Document!, record.DocumentId, serializerOptions: null);
SinkEnvelopeJson.WriteMetadata(writer, record.Metadata); // a "metadata" property
writer.WriteEndObject();The envelope shape around these pieces is yours to define.
SinkDestination.Resolve(record, sinkDefault, Name, "DefaultIndex") is the destination fallback the built-in sinks share: the record's destination, else the sink's default, else a WallabyConfigurationException naming the option to set.
The delegate sink
For in-process handlers (tests, side-effects, quick integrations), you can skip the class and use a lambda:
cdc.AddDelegateSink("audit", async (batch, ct) =>
{
foreach (var r in batch.Records)
{
if (r.IsDeletion) await store.RemoveAsync(r.DocumentId, ct);
else await store.UpsertAsync(r.DocumentId, r.Document!, ct);
}
return DeliveryResult.Success;
});Example
TIP
A production-ready HTTP sink ships as Wallaby.Sinks.Http with batched envelopes, HMAC signing, and IHttpClientFactory integration. The example below is a simple example.
public sealed class HttpSink(HttpClient http) : ISink
{
public string Name => "http";
public async Task<DeliveryResult> DeliverAsync(SinkBatch batch, CancellationToken ct)
{
try
{
foreach (var r in batch.Records)
{
using var resp = r.IsDeletion
? await http.DeleteAsync($"/docs/{r.DocumentId}", ct)
: await http.PutAsJsonAsync($"/docs/{r.DocumentId}", r.Document, ct);
if ((int)resp.StatusCode >= 500) return DeliveryResult.Retry($"upstream {(int)resp.StatusCode}");
if (!resp.IsSuccessStatusCode) return DeliveryResult.Permanent($"rejected {(int)resp.StatusCode}");
}
return DeliveryResult.Success;
}
catch (HttpRequestException ex) { return DeliveryResult.Retry(ex.Message, ex); }
}
}