Loading

Channel hierarchy

IngestChannel<T> is the top of a layered class hierarchy. Each layer adds a specific capability.

BufferedChannelBase~TOptions,TEvent,TResponse~ + TryWrite(TEvent): bool + WaitToWriteAsync(TEvent, CancellationToken): Task~bool~ + TryWriteMany(IEnumerable~TEvent~): bool + WaitForDrainAsync(TimeSpan?, CancellationToken): Task # ExportAsync(ArraySegment~TEvent~, CancellationToken): Task~TResponse~ # RetryBuffer(TResponse, list, buffer): list ResponseItemsBufferedChannelBase~...,TBulkResponseItem~ # Retry(TResponse): bool # RetryAllItems(TResponse): bool # RetryEvent(TEvent, TBulkResponseItem): bool # RejectEvent(TEvent, TBulkResponseItem): bool # Zip(TResponse, list): IEnumerable TransportChannelBase~...,TBulkResponseItem~ # ExportAsync(ITransport, ArraySegment~TEvent~, CancellationToken): Task~TResponse~ IngestChannelBase~TDocument,TOptions~ # ExportAsync(ITransport, ArraySegment, CancellationToken): Task~BulkResponse~ # Retry(BulkResponse): bool # RetryAllItems(BulkResponse): bool # RetryEvent(TDocument, BulkResponseItem): bool # RejectEvent(TDocument, BulkResponseItem): bool IngestChannel<TEvent> + string: TemplateName + string: TemplateWildcard + IIndexProvisioningStrategy: ProvisioningStrategy + BootstrapElasticsearchAsync(): Task~bool~ + ApplyAliasesAsync(string, CancellationToken): Task~bool~ + RolloverAsync(string?, string?, long?, CancellationToken): Task~bool~ BufferedChannelBase ResponseItemsBufferedChannelBase TransportChannelBase IngestChannelBase

The foundation layer in Elastic.Channels. Provides:

  • Two-stage inbound/outbound buffer (see push model)
  • Concurrent export with throttling
  • Backpressure and drain
  • Retry with configurable backoff

All buffer configuration, threading, and lifecycle management lives here.

Adds per-item response analysis. Where the base class treats a response as pass/fail for the whole batch, this layer can:

  • Zip sent items with individual response items
  • Decide per-item whether to retry, reject, or accept
  • Retry the entire batch for certain response codes (for example, HTTP 429)

In Elastic.Ingest.Transport. Integrates ITransport from Elastic.Transport:

  • Passes the transport instance to ExportAsync
  • Subclasses implement ExportAsync(ITransport, ArraySegment<TEvent>, CancellationToken) to make HTTP calls

In Elastic.Ingest.Elasticsearch. Implements the Elasticsearch _bulk API:

  • Serializes documents into NDJSON bulk format
  • Parses BulkResponse and maps items
  • Defines retry status codes (429, 500-599)
  • Defines rejection criteria (non-2xx that aren't retryable)

The composable channel that users interact with. Adds:

  • Strategy pattern for bootstrap, ingest, provisioning, alias, and rollover
  • Auto-configuration from ElasticsearchTypeContext
  • Public API for bootstrap, alias management, and rollover

Before IngestChannel<T>, specialized channel subclasses existed for each use case. These are now obsolete. See legacy channels for migration guidance.