Bulk sender and pull ingestion

Channels are a push model: you write events, the channel owns the threads, batching, concurrency and retry. Two other shapes are common:

You have Use
One batch you built yourself, and your own retry/ordering/threading BulkSender: one request per call, no threads, nothing to dispose
A finite list or IAsyncEnumerable you want stored IngestAllAsync: pulls, batches, sends and returns when everything settled
A configured channel (strategies, bootstrap, callbacks) and a finite sequence channel.IngestAllAsync

None of these require Elastic.Mapping, an [Index<T>] attribute or a bootstrap step. All of them work with a plain ITransport and a JsonTypeInfo<T>, so they are AOT and trim safe when you use a source generated serializer context.

A BulkSender<TItem, TBody> is immutable and thread safe. Per call it serializes into a pooled buffer, sends one _bulk request, and returns the typed BulkResponse.

var sender = BulkSender.Create(transport, MyContext.Default.Product,
    action: static p => BulkAction.Index(id: p.Sku), target: "products");

BulkResponse response = await sender.SendAsync(products, ct);
if (!response.AllItemsPersisted()) { /* inspect response.Items */ }
		
  1. POST products/_bulk

SendAsync accepts a ReadOnlySpan<TItem> (arrays, CollectionsMarshal.AsSpan(list), slices) or an IEnumerable<TItem>.

The Action delegate runs once per item and returns a BulkAction. It is a small struct, so this does not allocate.

BulkAction Written as
Index(id?, index?) {"index":{...}} + document
Create(id?, index?) {"create":{...}} + document (409 if the id exists)
Update(id, index?) {"update":{...}} + {"doc_as_upsert":true,"doc":...}
Delete(id, index?) {"delete":{...}} only, no body line
ScriptedHashUpsert(id, HashedBulkUpdate, index?) scripted upsert that only replaces the document when its hash changed

.WithRequireAlias() and .WithDynamicTemplates(...) modify an action. When index is null the sender's Target applies ({target}/_bulk); when both are null the request goes to _bulk and the action must carry the index.

TItem and TBody are separate so the action can be derived from a change record while the body is just the payload:

var sender = new BulkSender<ChangeRecord, JsonElement>(new()
{
    Transport    = transport,
    BodyTypeInfo = MyContext.Default.JsonElement,
    Body         = static c => c.After,
    Action       = static c => c.Op switch
    {
        Op.Insert => BulkAction.Index(c.Key, c.IndexName),
        Op.Update => BulkAction.Update(c.Key, c.IndexName),
        Op.Delete => BulkAction.Delete(c.Key, c.IndexName),
        _ => throw new ArgumentOutOfRangeException()
    },
    Retry = BulkRetryPolicy.None
});
		
  1. not called for deletes
  2. the caller owns retry

The default is BulkRetryPolicy.None: exactly one request. BulkRetryPolicy.Default retries items that failed with 429, 502, 503 or 504 up to three times, two seconds apart. Everything is replaceable:

Retry = BulkRetryPolicy.Default with
{
    MaxRetries  = 5,
    Backoff     = static n => TimeSpan.FromMilliseconds(100 * (1 << n)),
    IsRetryable = static item => item.Status is 429 or 503
}
		

Only the failed items are re-sent, without serializing them again. response.Items always lines up position by position with the items you passed in, also after retries: each position holds the last result for that item, so response.Items[i] belongs to items[i]. A top level HTTP 429 re-sends the whole request (RetryAllOnHttp429).

Two options apply to every request a sender sends, including retries and the batches of IngestAllAsync:

var sender = BulkSender.Create(transport, MyContext.Default.Order,
    action: static o => BulkAction.Index(id: o.Id),
    target: "orders",
    refresh: BulkRefresh.WaitFor,
    requestTimeout: TimeSpan.FromSeconds(30));
		
  1. POST orders/_bulk?refresh=wait_for&filter_path=...
Option Effect
Refresh Sets the refresh parameter: BulkRefresh.WaitFor returns once the documents are visible to search without forcing a refresh, True refreshes the affected shards before returning, False sends refresh=false. null (the default) sends nothing and Elasticsearch refreshes on the index's own interval.
RequestTimeout The client side timeout of each _bulk request, layered on top of the transport's own configuration. null (the default) leaves the transport's timeout in place. It must be positive or Timeout.InfiniteTimeSpan.

RequestTimeout is the client's wait for the HTTP response. It is not Elasticsearch's server side timeout parameter, which limits how long the cluster waits for active shards or mapping updates.

Both are fixed for the lifetime of a sender. To vary them per call, keep one sender per setting. A sender is cheap, immutable and thread safe. There is deliberately no free form query string, so these options cannot collide with the built in filter_path.

BulkResponseItem.Id and BulkResponseItem.Index report the _id and _index Elasticsearch used for each item. They are the only way to learn an id the server generated (BulkAction.Index() or Create() without an id) or the concrete index behind an alias (WithRequireAlias()) or a data stream. Neither is reported by default; ask for the fields you use with ReturnItemIdentity:

var sender = new BulkSender<Event, Event>(new BulkSenderOptions<Event, Event>
{
    /* transport, action, body ... */
    ReturnItemIdentity = Track.Id
});

BulkResponse response = await sender.SendAsync(events, ct);
var generatedId = response.Items.First().Id;              // "AXx9kQ3pTzGm..." assigned by Elasticsearch
		
  1. or Track.Index, or Track.Id | Track.Index

BulkSender.Create(..., itemIdentity: Track.Id | Track.Index) takes the same flags, and the channel options have the same ReturnItemIdentity property.

Flag Reports Use it for
Track.Id Id: the generated id, or the explicit id echoed back learning ids Elasticsearch assigned
Track.Index Index: the concrete index that received the item learning the backing index behind an alias or data stream

Each field costs something, so request only what you use. A response item that only carries an action and a status is a shared, immutable instance and allocates nothing. An item that carries an _id is its own object (about 56 bytes) plus the id string (about 64 bytes), and the response is larger because every item repeats the field. The index string is reused across the items of a response, so Track.Index costs about half of Track.Id. Measured for 1,000 items: 18 KB without either field, 74 KB with Track.Index, 138 KB with Track.Id or both. Many workloads, for example logs written to a data stream, need neither and pay nothing.

The flags apply to every request the sender issues, retries and the batches of IngestAllAsync included. response.Items[i] still belongs to items[i], carrying the identity from the attempt that settled it, and IngestAllAsync failures expose the same properties. A request that failed as a whole has no items, so the failures made up for it have Id and Index set to null.

Channels work the same way: set ReturnItemIdentity = Track.Id | Track.Index on the channel options to have DirectWriteAsync, the response callbacks and IngestAllAsync see Id and Index for every item.

By default a sender serializes documents the same way channels do: with DefaultIgnoreCondition = WhenWritingDefault (members holding their default value, such as int N = 0 or a null string, are omitted). Output from BulkSender and from a channel is therefore identical for the same document.

This default is applied on top of your JsonTypeInfo: it overrides a DefaultIgnoreCondition set on your serializer context (via [JsonSourceGenerationOptions]), exactly as it does for channels. Per property attributes such as [JsonIgnore(Condition = ...)], custom converters and naming policies keep working. To serialize exactly as your type info is configured, set ApplyLibrarySerializerDefaults = false on BulkSenderOptions.

For "I have a list, store it":

BulkIngestAllResult result = await BulkSender.IngestAllAsync(
    transport, MyContext.Default.Product, products, target: "products");

if (result.Failures.Count > 0)
    foreach (var f in result.Failures)
        Console.WriteLine($"{f.Position}: {f.Item.Status} {f.Item.Error?.Reason}");
		

It batches the source (BatchSize, default 1,000), sends up to MaxConcurrency requests at a time (default: processor count, 1 keeps strict order), retries with BulkRetryPolicy.Default unless IngestAllOptions.Retry says otherwise, and returns when every batch settled. Failures are the items that were still failing after retries, with their position in the source. Arrays and lists are sliced without copying; other sequences are copied into pooled arrays. There is an IAsyncEnumerable<T> overload, and the same methods exist as instance methods on a configured BulkSender.

When you want a channel's strategies, bootstrap and callbacks but already hold a finite sequence:

IngestAllResult result = await channel.IngestAllAsync(products, maxConcurrency: 1, ctx);
		

The channel fills batches of BatchExportSize straight from the source and exports them with the same retry and callback machinery used for pushed events. It bypasses the inbound buffer: no threads are started, InflightEvents is untouched, the channel is not completed, and the task completes exactly when every batch settled, so no WaitForDrainAsync is needed. The channel stays usable: call IngestAllAsync repeatedly, or mix it with TryWrite.

  • maxConcurrency defaults to the channel's MaxConcurrency; 1 preserves order.
  • For IAsyncEnumerable sources a partial batch older than OutboundBufferMaxLifetime is exported even while the source is stalled.
  • IngestAllResult.RetriesExhausted counts items that were still failing when ExportMaxRetries ran out. Permanently rejected items are reported through ServerRejectionCallback as for pushed events.
  • This works for a channel built from plain IndexChannelOptions<T> (an IndexFormat and a transport) as well as for the [Index<T>] mapping context path.