Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion quickwit/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion quickwit/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ metrics-util = "0.20"
mime_guess = "2.0"
mini-moka = "0.10"
mockall = "0.14"
mrecordlog = { git = "https://github.com/quickwit-oss/mrecordlog", rev = "3b3562ef" }
mrecordlog = { git = "https://github.com/quickwit-oss/mrecordlog", rev = "582ddb131181a053635e6904d3190491a3ee556c" }
new_string_template = "1.5"
nom = "8.0"
numfmt = "1.2"
Expand Down
8 changes: 4 additions & 4 deletions quickwit/quickwit-indexing/src/source/ingest/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1560,7 +1560,7 @@ mod tests {
from_position_exclusive: Some(Position::offset(11u64)),
to_position_inclusive: Some(Position::offset(14u64)),
};
let batch_size = fetch_payload.estimate_size();
let batch_size = fetch_payload.buffer_size();
let fetch_message = FetchMessage::new_payload(fetch_payload);
let in_flight_value =
InFlightValue::new(fetch_message, batch_size, &IN_FLIGHT_FETCH_STREAM);
Expand All @@ -1574,7 +1574,7 @@ mod tests {
from_position_exclusive: Some(Position::offset(22u64)),
to_position_inclusive: Some(Position::offset(23u64)),
};
let batch_size = fetch_payload.estimate_size();
let batch_size = fetch_payload.buffer_size();
let fetch_message = FetchMessage::new_payload(fetch_payload);
let in_flight_value =
InFlightValue::new(fetch_message, batch_size, &IN_FLIGHT_FETCH_STREAM);
Expand Down Expand Up @@ -1653,7 +1653,7 @@ mod tests {
from_position_exclusive: Some(Position::offset(14u64)),
to_position_inclusive: Some(Position::offset(15u64)),
};
let batch_size = fetch_payload.estimate_size();
let batch_size = fetch_payload.buffer_size();
let fetch_message = FetchMessage::new_payload(fetch_payload);
let in_flight_value =
InFlightValue::new(fetch_message, batch_size, &IN_FLIGHT_FETCH_STREAM);
Expand Down Expand Up @@ -1795,7 +1795,7 @@ mod tests {
from_position_exclusive: Some(Position::offset(11u64)),
to_position_inclusive: Some(Position::offset(13u64)),
};
let batch_size = fetch_payload.estimate_size();
let batch_size = fetch_payload.buffer_size();
let fetch_message = FetchMessage::new_payload(fetch_payload);
let in_flight_value =
InFlightValue::new(fetch_message, batch_size, &IN_FLIGHT_FETCH_STREAM);
Expand Down
2 changes: 1 addition & 1 deletion quickwit/quickwit-ingest/src/ingest_v2/doc_mapper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ fn is_document_validation_enabled() -> bool {
!quickwit_common::get_bool_from_env_cached!("QW_DISABLE_DOCUMENT_VALIDATION", false)
}

#[instrument(name = "ingester.validate_doc_batch", skip_all, fields(num_docs = doc_batch.num_docs(), num_bytes = doc_batch.num_bytes()))]
#[instrument(name = "ingester.validate_doc_batch", skip_all, fields(num_docs = doc_batch.num_docs(), num_bytes = doc_batch.doc_buffer.len()))]
async fn validate_doc_batch_cpu_intensive(
doc_batch: DocBatchV2,
doc_mapper: Arc<DocMapper>,
Expand Down
4 changes: 2 additions & 2 deletions quickwit/quickwit-ingest/src/ingest_v2/fetch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -170,7 +170,7 @@ impl FetchStreamTask {
mrecord_buffer: mrecord_buffer.freeze(),
mrecord_lengths,
};
let batch_size = mrecord_batch.estimate_size();
let batch_size = mrecord_batch.buffer_size();
let fetch_payload = FetchPayload {
index_uid: Some(self.index_uid.clone()),
source_id: self.source_id.clone(),
Expand Down Expand Up @@ -490,7 +490,7 @@ async fn fetch_stream_once(
match fetch_message_result {
Ok(fetch_message) => match &fetch_message.message {
Some(fetch_message::Message::Payload(fetch_payload)) => {
let batch_size = fetch_payload.estimate_size();
let batch_size = fetch_payload.buffer_size();
let to_position_inclusive = fetch_payload.to_position_inclusive();
let in_flight_value = InFlightValue::new(
fetch_message,
Expand Down
118 changes: 105 additions & 13 deletions quickwit/quickwit-ingest/src/ingest_v2/ingester.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,11 +53,11 @@ use super::local_shards::ShardThroughputReadings;
use super::metrics::report_local_shards_metrics;
use super::models::IngesterShard;
use super::mrecordlog_utils::{
AppendDocBatchError, append_non_empty_doc_batch, check_enough_capacity, wal_stats,
AppendDocBatchError, append_non_empty_doc_batch, check_enough_capacity, doc_batch_size,
wal_stats,
};
use super::rate_meter::RateMeter;
use super::state::{IngesterState, InnerIngesterState, WeakIngesterState};
use crate::estimate_size;
use crate::ingest_v2::doc_mapper::get_or_try_build_doc_mapper;
use crate::ingest_v2::metrics::{
DECOMMISSION_FAILED, DECOMMISSION_SUCCEEDED, RESET_SHARDS_OPERATIONS_TOTAL, STATUS,
Expand Down Expand Up @@ -476,7 +476,7 @@ impl Ingester {
DocBatchV2::default()
}
};
let requested_capacity = estimate_size(&doc_batch);
let requested_capacity = doc_batch_size(&doc_batch, force_commit);

if let Err(error) = check_enough_capacity(
&state_guard.mrecordlog,
Expand Down Expand Up @@ -517,7 +517,7 @@ impl Ingester {
}

// Total number of bytes (valid and invalid documents)
let original_batch_num_bytes = doc_batch.num_bytes() as u64;
let original_batch_num_bytes = doc_batch.doc_buffer.len() as u64;

let (valid_doc_batch, parse_failures) = if validate_docs {
validate_doc_batch(doc_batch, doc_mapper).await?
Expand Down Expand Up @@ -558,7 +558,7 @@ impl Ingester {
parent: DOCS_BYTES_TOTAL,
labels: [label_values!(VALIDITY => "valid")],
)
.inc_by(valid_doc_batch.num_bytes() as u64);
.inc_by(valid_doc_batch.doc_buffer.len() as u64);
if !parse_failures.is_empty() {
counter!(
parent: DOCS_TOTAL,
Expand All @@ -569,9 +569,9 @@ impl Ingester {
parent: DOCS_BYTES_TOTAL,
labels: [label_values!(VALIDITY => "invalid")],
)
.inc_by(original_batch_num_bytes - valid_doc_batch.num_bytes() as u64);
.inc_by(original_batch_num_bytes - valid_doc_batch.doc_buffer.len() as u64);
}
let valid_batch_num_bytes = valid_doc_batch.num_bytes() as u64;
let valid_batch_num_bytes = valid_doc_batch.doc_buffer.len() as u64;
shard.rate_meter.update(valid_batch_num_bytes);
total_requested_capacity += requested_capacity;

Expand Down Expand Up @@ -606,8 +606,8 @@ impl Ingester {
)
.await;

let current_position_inclusive = match append_result {
Ok(current_position_inclusive) => current_position_inclusive,
let (current_position_inclusive, queue_size) = match append_result {
Ok(outcome) => outcome,
Err(append_error) => {
let reason = match &append_error {
AppendDocBatchError::Io(io_error) => {
Expand Down Expand Up @@ -637,11 +637,12 @@ impl Ingester {
}
};

state_guard
let shard = state_guard
.shards
.get_mut(&queue_id)
.expect("shard should exist")
.set_replication_position_inclusive(current_position_inclusive.clone(), now);
.expect("shard should exist");
shard.set_replication_position_inclusive(current_position_inclusive.clone(), now);
shard.queue_size = queue_size;

let persist_success = PersistSuccess {
subrequest_id: subrequest.subrequest_id,
Expand Down Expand Up @@ -925,7 +926,9 @@ impl IngesterService for Ingester {
.subrequests
.iter()
.flat_map(|subrequest| match &subrequest.doc_batch {
Some(doc_batch) if doc_batch.doc_buffer.is_unique() => Some(doc_batch.num_bytes()),
Some(doc_batch) if doc_batch.doc_buffer.is_unique() => {
Some(doc_batch.doc_buffer.len())
}
_ => None,
})
.sum::<usize>();
Expand Down Expand Up @@ -1683,6 +1686,87 @@ mod tests {
assert!(shard.long_term_ingestion_rate.as_u64() > 0);
}

#[tokio::test]
async fn test_queue_size_does_not_drift_after_persist_and_truncate() {
let (ingester_ctx, ingester) = IngesterForTest::default().build().await;
let index_uid = IndexUid::for_test("test-index", 0);
let source_id = SourceId::from("test-source");
let doc_mapping_uid = DocMappingUid::random();
let response = ingester
.init_shards(InitShardsRequest {
subrequests: vec![InitShardSubrequest {
subrequest_id: 0,
shard: Some(Shard {
index_uid: Some(index_uid.clone()),
source_id: source_id.clone(),
shard_id: Some(ShardId::from(1)),
shard_state: ShardState::Open as i32,
ingester_id: ingester_ctx.node_id.to_string(),
doc_mapping_uid: Some(doc_mapping_uid),
..Default::default()
}),
doc_mapping_json: format!(r#"{{"doc_mapping_uid":"{doc_mapping_uid}"}}"#),
validate_docs: false,
}],
})
.await
.unwrap();
assert_eq!(response.successes.len(), 1);
assert!(response.failures.is_empty());

let queue_id = queue_id(&index_uid, &source_id, &ShardId::from(1));
let mut next_position = 0u64;
let mut retained_bytes = 0u64;
for iteration in 0..32 {
let force_commit = iteration % 2 == 0;
let commit_type = if force_commit {
CommitTypeV2::Force
} else {
CommitTypeV2::Auto
};
let response = ingester
.persist(PersistRequest {
ingester_id: ingester_ctx.node_id.to_string(),
commit_type: commit_type as i32,
subrequests: vec![PersistSubrequest {
subrequest_id: 0,
index_uid: Some(index_uid.clone()),
source_id: source_id.clone(),
doc_batch: Some(DocBatchV2::for_test(["a", "longer", "last"])),
}],
})
.await
.unwrap();
assert_eq!(response.successes.len(), 1);
assert!(response.failures.is_empty());

let mut state_guard = ingester.state.lock_fully("test").await.unwrap();
let commit_bytes = u64::from(force_commit) * 2;
assert_eq!(
state_guard.shards[&queue_id].queue_size,
ByteSize::b(retained_bytes + 17 + commit_bytes)
);
state_guard
.truncate_shard(&queue_id, Position::offset(next_position + 1), "test")
.await;
retained_bytes = 6 + commit_bytes;
assert_eq!(
state_guard.mrecordlog.summary().queues[&queue_id].num_bytes as u64,
retained_bytes
);
assert_eq!(
state_guard.shards[&queue_id].queue_size,
ByteSize::b(retained_bytes)
);
next_position += 3 + u64::from(force_commit);
}
let mut state_guard = ingester.state.lock_fully("test").await.unwrap();
state_guard
.truncate_shard(&queue_id, Position::offset(next_position - 1), "test")
.await;
assert_eq!(state_guard.shards[&queue_id].queue_size, ByteSize::b(0));
}

#[tokio::test]
async fn test_ingester_persist() {
let (ingester_ctx, ingester) = IngesterForTest::default().build().await;
Expand Down Expand Up @@ -1782,6 +1866,10 @@ mod tests {
let shard_01 = state_guard.shards.get(&queue_id_01).unwrap();
shard_01.assert_is_open();
shard_01.assert_replication_position(Position::offset(1u64));
assert_eq!(
shard_01.queue_size,
ByteSize::b(state_guard.mrecordlog.summary().queues[&queue_id_01].num_bytes as u64)
);

state_guard.mrecordlog.assert_records_eq(
&queue_id_01,
Expand All @@ -1793,6 +1881,10 @@ mod tests {
let shard_11 = state_guard.shards.get(&queue_id_11).unwrap();
shard_11.assert_is_open();
shard_11.assert_replication_position(Position::offset(2u64));
assert_eq!(
shard_11.queue_size,
ByteSize::b(state_guard.mrecordlog.summary().queues[&queue_id_11].num_bytes as u64)
);

state_guard.mrecordlog.assert_records_eq(
&queue_id_11,
Expand Down
35 changes: 7 additions & 28 deletions quickwit/quickwit-ingest/src/ingest_v2/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,6 @@ use workbench::pending_subrequests;

pub use self::fetch::{FetchStreamError, MultiFetchStream};
pub use self::ingester::Ingester;
use self::mrecord::MRECORD_HEADER_LEN;
pub use self::mrecord::{MRecord, decoded_mrecords};
pub use self::router::IngestRouter;

Expand Down Expand Up @@ -280,11 +279,6 @@ impl IngestRequestV2Builder {
}
}

pub(super) fn estimate_size(doc_batch: &DocBatchV2) -> ByteSize {
let estimate = doc_batch.num_bytes() + doc_batch.num_docs() * MRECORD_HEADER_LEN;
ByteSize(estimate as u64)
}

#[derive(Debug, Clone, Copy, Default, Eq, PartialEq, Ord, PartialOrd)]
pub struct RateMibPerSec(pub u16);

Expand Down Expand Up @@ -335,7 +329,7 @@ mod tests {
let doc_batch = doc_batch_builder.build().unwrap();

assert_eq!(doc_batch.num_docs(), 2);
assert_eq!(doc_batch.num_bytes(), 21);
assert_eq!(doc_batch.doc_buffer.len(), 13);
assert_eq!(doc_batch.doc_lengths, [7, 6]);
assert_eq!(doc_batch.doc_buffer, Bytes::from(&b"Hello, World!"[..]));
}
Expand Down Expand Up @@ -391,8 +385,9 @@ mod tests {
.doc_batch
.as_ref()
.unwrap()
.num_bytes(),
21
.doc_buffer
.len(),
13
);
assert_eq!(
ingest_request.subrequests[0]
Expand Down Expand Up @@ -434,8 +429,9 @@ mod tests {
.doc_batch
.as_ref()
.unwrap()
.num_bytes(),
20
.doc_buffer
.len(),
12
);
assert_eq!(
ingest_request.subrequests[1]
Expand All @@ -462,21 +458,4 @@ mod tests {
[hola_doc_uid, mundo_doc_uid]
);
}

#[test]
fn test_estimate_size() {
let doc_batch = DocBatchV2 {
doc_buffer: Vec::new().into(),
doc_lengths: Vec::new(),
doc_uids: Vec::new(),
};
assert_eq!(estimate_size(&doc_batch), ByteSize(0));

let doc_batch = DocBatchV2 {
doc_buffer: vec![0u8; 100].into(),
doc_lengths: vec![10, 20, 30],
doc_uids: Vec::new(),
};
assert_eq!(estimate_size(&doc_batch), ByteSize(118));
}
}
Loading
Loading