From 40a4629f2d0239a735172027164814a98a144349 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Thu, 20 Aug 2026 22:55:33 +0200 Subject: [PATCH] Add create-only split staging mode Expose create-only staging semantics in the metastore API and use them during reconciliation so existing split rows are preserved across retries and concurrent writers. --- .../file_backed/file_backed_index/mod.rs | 11 ++++++ .../src/metastore/file_backed/mod.rs | 8 +++- .../quickwit-metastore/src/metastore/mod.rs | 2 + .../src/metastore/postgres/metastore.rs | 8 +++- .../quickwit-metastore/src/tests/split.rs | 37 +++++++++++++++++++ .../protos/quickwit/metastore.proto | 2 + .../codegen/quickwit/quickwit.metastore.rs | 3 ++ 7 files changed, 68 insertions(+), 3 deletions(-) diff --git a/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/mod.rs b/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/mod.rs index ba1d226fc74..b149546693e 100644 --- a/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/mod.rs +++ b/quickwit/quickwit-metastore/src/metastore/file_backed/file_backed_index/mod.rs @@ -333,6 +333,17 @@ impl FileBackedIndex { Ok(()) } + /// Stages a split only when its ID is absent, preserving any concurrent writer's row. + pub(crate) fn stage_split_create_only( + &mut self, + split_metadata: SplitMetadata, + ) -> Result<(), MetastoreError> { + if self.splits.contains_key(split_metadata.split_id()) { + return Ok(()); + } + self.stage_split(split_metadata) + } + /// Marks the splits for deletion. Returns whether a mutation occurred. pub(crate) fn mark_splits_for_deletion( &mut self, diff --git a/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs b/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs index 20e573ac139..80be384ccb9 100644 --- a/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs +++ b/quickwit/quickwit-metastore/src/metastore/file_backed/mod.rs @@ -656,13 +656,19 @@ impl MetastoreService for FileBackedMetastore { #[instrument(name = "metastore.file_backed.stage_splits", skip_all, fields(index_uid = %request.index_uid()))] async fn stage_splits(&self, request: StageSplitsRequest) -> MetastoreResult { let index_uid = request.index_uid().clone(); + let create_only = request.create_only; let splits_metadata = request.deserialize_splits_metadata()?; self.mutate(&index_uid, |index| { let mut failed_split_ids = Vec::new(); for split_metadata in splits_metadata { - match index.stage_split(split_metadata) { + let stage_result = if create_only { + index.stage_split_create_only(split_metadata) + } else { + index.stage_split(split_metadata) + }; + match stage_result { Ok(()) => {} Err(MetastoreError::FailedPrecondition { entity: EntityKind::Split { split_id }, diff --git a/quickwit/quickwit-metastore/src/metastore/mod.rs b/quickwit/quickwit-metastore/src/metastore/mod.rs index c12b8b11b61..a5420356b62 100644 --- a/quickwit/quickwit-metastore/src/metastore/mod.rs +++ b/quickwit/quickwit-metastore/src/metastore/mod.rs @@ -650,6 +650,7 @@ impl StageSplitsRequestExt for StageSplitsRequest { let request = Self { index_uid: Some(index_uid.into()), split_metadata_list_serialized_json, + create_only: false, }; Ok(request) } @@ -663,6 +664,7 @@ impl StageSplitsRequestExt for StageSplitsRequest { let request = Self { index_uid: Some(index_uid.into()), split_metadata_list_serialized_json, + create_only: false, }; Ok(request) } diff --git a/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs b/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs index 2bf2fec652f..736a66e989e 100644 --- a/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs +++ b/quickwit/quickwit-metastore/src/metastore/postgres/metastore.rs @@ -679,6 +679,7 @@ impl MetastoreService for PostgresqlMetastore { #[instrument(name = "metastore.postgres.stage_splits", skip_all, fields(split_ids))] async fn stage_splits(&self, request: StageSplitsRequest) -> MetastoreResult { let index_uid: IndexUid = request.index_uid().clone(); + let create_only = request.create_only; let splits_metadata = request.deserialize_splits_metadata()?; if splits_metadata.is_empty() { @@ -745,7 +746,9 @@ impl MetastoreService for PostgresqlMetastore { node_id = excluded.node_id, update_timestamp = CURRENT_TIMESTAMP, create_timestamp = CURRENT_TIMESTAMP - WHERE splits.split_id = excluded.split_id AND splits.split_state = 'Staged' + WHERE splits.split_id = excluded.split_id + AND splits.split_state = 'Staged' + AND NOT $11 RETURNING split_id; "#) .bind(&split_ids) @@ -758,11 +761,12 @@ impl MetastoreService for PostgresqlMetastore { .bind(&node_ids) .bind(SplitState::Staged.as_str()) .bind(&index_uid) + .bind(create_only) .fetch_all(tx.as_mut()) .await .map_err(|sqlx_error| convert_sqlx_err(&index_uid.index_id, sqlx_error))?; - if upserted_split_ids.len() != split_ids.len() { + if !create_only && upserted_split_ids.len() != split_ids.len() { let failed_split_ids: Vec = split_ids .into_iter() .filter(|split_id| !upserted_split_ids.contains(split_id)) diff --git a/quickwit/quickwit-metastore/src/tests/split.rs b/quickwit/quickwit-metastore/src/tests/split.rs index a66e91a84c4..6a8b7ff6cdb 100644 --- a/quickwit/quickwit-metastore/src/tests/split.rs +++ b/quickwit/quickwit-metastore/src/tests/split.rs @@ -1621,6 +1621,43 @@ pub async fn test_metastore_stage_splits, #[prost(string, tag = "2")] pub split_metadata_list_serialized_json: ::prost::alloc::string::String, + /// Create-only mode: insert missing rows without upserting an existing split row. + #[prost(bool, tag = "3")] + pub create_only: bool, } #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]