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 :)
Problem
The
elasticsearch_v8/elasticsearch_v9outputs can set a document_idand anaction, but cannot set Elasticsearch's per-document external version. Each bulk op is built with onlyIndex_,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_idarrive in order, but it's less well suited to streaming ingest where the same_idis written repeatedly and correctness depends on which version wins, not which arrived last:max_in_flight > 1(concurrent bulks, per-shard fan-out) can apply two writes to the same_idin the wrong order.Proposal
Add an opt-in, interpolated
versionfield. 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 returns409, 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
With this config, replaying an older record for an existing
idis a no-op (Elasticsearch returns409, which the output treats as success) instead of overwriting the newer document.Open items / possible follow-ups
external_gte: This iteration scopes toversion_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 :)