OpenSearch Sink
Keep OpenSearch indexes in sync with Postgres from your .NET application. The Wallaby.Sinks.OpenSearch package streams committed row changes out of Postgres logical replication and delivers them through the _bulk API: upserts are indexed with _id set to a stable document id (so updates are idempotent) and deletions remove by that same id. No polling, no dual writes, no reindex script. It works with self-managed OpenSearch and Amazon OpenSearch Service.
Quickstart
dotnet add package Wallaby.Sinks.OpenSearchRegister Wallaby, point it at a storage provider, add the sink, and map an entity. The mapping's destination is the index name (index names must be lowercase).
builder.Services.AddWallaby(cdc =>
{
cdc.UseEntityFrameworkCore<AppDbContext>()
.UseConnectionString(conn)
.AddOpenSearchSink("search", s =>
{
s.Endpoint = "https://localhost:9200";
s.Username = "wallaby";
s.Password = password;
s.DefaultIndex = "documents";
})
.WithMappings(sink => sink
.Map<Product>()
.ToDestination("products")
.UsingTransform(/* ... */));
});builder.Services.AddWallaby(cdc =>
{
cdc.UseMarten()
.UseConnectionString(conn)
.AddOpenSearchSink("search", s =>
{
s.Endpoint = "https://localhost:9200";
s.Username = "wallaby";
s.Password = password;
s.DefaultIndex = "documents";
})
.WithMappings(sink => sink
.Map<Product>()
.ToDestination("products")
.UsingTransform(/* ... */));
});builder.Services.AddWallaby(cdc =>
{
cdc.UseTables(tables => tables.Add<Product>())
.UseConnectionString(conn)
.AddOpenSearchSink("search", s =>
{
s.Endpoint = "https://localhost:9200";
s.Username = "wallaby";
s.Password = password;
s.DefaultIndex = "documents";
})
.WithMappings(sink => sink
.Map<Product>()
.ToDestination("products")
.UsingTransform(/* ... */));
});The transform shapes each change into the document you want indexed; see mappings. For the Postgres server settings Wallaby needs, see getting started.
Options
| Option | Default | Purpose |
|---|---|---|
Endpoint | (required) | OpenSearch base URL. |
Username / Password | null | Basic auth; null for unsecured. |
ConfigureConnection | null | Full override for the client's connection settings (see below). |
DefaultIndex | null | Index used when a routed record has no destination; a record with neither fails permanently. |
MaxRecordsPerRequest | 500 | Records per _bulk request; larger batches are split into sequential requests, preserving commit order. |
Timeout | 30s | Per-request timeout. |
Refresh | false | When true, bulk requests use refresh=wait_for so documents are searchable before the batch is acknowledged. |
SerializerOptions | null | Serializer for document values beyond the natively written scalar types (numbers, strings, dates, byte[], vectors). |
Indexes
The sink doesn't create or configure indexes. An index is created automatically on first write with dynamic mapping, as long as the cluster's action.auto_create_index setting allows it (it does by default). For explicit settings or mappings (analyzers, knn_vector fields, shard counts, …), create the index up front with Dev Tools, your infrastructure tooling, or a deployment script. In-sink index bootstrapping is planned.
Vector search
A transform can emit an embedding as a float[] or ReadOnlyMemory<float> field. The sink writes either as a plain JSON number array, so no SerializerOptions are needed. Dynamic mapping would infer an ordinary float field, so create the index up front with an explicit knn_vector mapping:
PUT /products
{
"settings": { "index.knn": true },
"mappings": {
"properties": {
"embedding": { "type": "knn_vector", "dimension": 1536 }
}
}
}Don't pass a quantized vector as byte[]: byte arrays serialize as base64 strings, not arrays.
OpenSearch can also do the embedding for you. Attach a neural-search ingest pipeline (a text_embedding processor over a deployed model) to the index and sync plain text, and the cluster computes vectors at index time, with none in your pipeline at all. The pipeline runs on every indexed document, so live changes embed incrementally, but a backfill re-runs the model over the whole corpus. See RAG & Embeddings.
Authentication
Username/Password cover basic auth. For anything else (AWS SigV4, client certificates, connection pools, proxies), take over construction of the client's connection settings with ConfigureConnection:
// Amazon OpenSearch Service, signed with the host's AWS credentials
// (requires the OpenSearch.Net.Auth.AwsSigV4 package):
cdc.AddOpenSearchSink("search", s =>
{
s.Endpoint = "https://my-domain.eu-west-1.es.amazonaws.com";
s.ConfigureConnection = uri => new ConnectionSettings(
new SingleNodeConnectionPool(uri), new AwsSigV4HttpConnection(RegionEndpoint.EUWest1));
});When ConfigureConnection is set, leave Username and Password unset (registration fails otherwise) and configure all authentication on the returned settings. Timeout still applies per request.
Purging
The sink implements purge-then-backfill: a purge runs _delete_by_query with match_all against the mapping's index (conflicts=proceed, refresh=true). It runs synchronously under the per-request Timeout, so a very large index may need a longer timeout. If the index doesn't exist yet, there's nothing to purge.
Delivery semantics
Delivery is at-least-once: a batch may be re-sent after a transient failure, and every action is idempotent by _id, so replays converge to the same documents. Batches are chunked into sequential _bulk requests so commit order is preserved.
Failures are classified per response and per bulk item:
- Throttling and server errors (
408/429/5xx, connection failures, timeouts) are retryable: the dispatcher backs off and re-sends. - Request or item rejections (e.g.
mapper_parsing_exceptionfrom a mapping conflict) are permanent. They point to a bug in a transform or the configuration, so the pipeline halts rather than silently dropping documents. - Deleting a document that's already gone reports
404for that item, which the sink treats as success.
By default, documents become searchable on the index's refresh interval (typically 1s) after the batch is acknowledged. Set Refresh = true to make each batch searchable before it's acknowledged, at a cost to indexing throughput.