Skip to content

elasticsearch output: support external document versioning #4893

Description

@yaakovkru

Problem

The elasticsearch_v8 / elasticsearch_v9 outputs can set a document _id and an action, but cannot set Elasticsearch's per-document external version. Each bulk op is built with only Index_, Id_, Pipeline, Routing, so the outputs always use ES internal versioning, hence last-writer-by-arrival. Internal versioning works well when writes to a given _id arrive in order, but it's less well suited to streaming ingest where the same _id is written repeatedly and correctness depends on which version wins, not which arrived last:

  • Out-of-order delivery with max_in_flight > 1 (concurrent bulks, per-shard fan-out) can apply two writes to the same _id in the wrong order.
  • At-least-once redelivery on replay an older offset document over a newer one.
  • Concurrent writers race, and internal versioning can't arbitrate by source order.

Proposal

Add an opt-in, interpolated version field. When set, the output assigns an external version (version_type: external, strict greater-than) to whole-document write actions (index, upsert, create, delete). ES applies the write only if the supplied version beats the stored one; otherwise it returns 409, which the output treats as a successful no-op making ingest idempotent and order-independent at the datastore. Fully opt-in (unset ⇒ remains current behavior).
Not applicable to update (update API uses internal versioning + retry_on_conflict).

Example

output:
  elasticsearch_v8:
    urls: ['http://localhost:9200']
    index: "things"
    action: "index"
    id: ${! metadata("kafka_key") }
    # Use the source record's timestamp as the version, so a later revision of
    # the same key always wins. Any monotonic per-key value works — a value from
    # a record header does too, e.g.
    #   version: ${! metadata("event_version") }
    version: ${! metadata("kafka_timestamp_ms") }

With this config, replaying an older record for an existing id is a no-op (Elasticsearch returns 409, which the output treats as success) instead of overwriting the newer document.

Open items / possible follow-ups

  • external_gte: This iteration scopes to version_type: external (strict >). external_gte (allow equal-version rewrites) could be added later if there's demand.

Implementation is ready; happy to open the PR :)

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions