Skip to content
Closed
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 @@ -181,7 +181,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 = "e79bed2dbe4077f02479d073574c9f7a7ab03b85" }
new_string_template = "1.5"
nom = "8.0"
numfmt = "1.2"
Expand Down
20 changes: 16 additions & 4 deletions quickwit/quickwit-ingest/src/ingest_v2/ingester.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,8 @@ 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};
Expand Down Expand Up @@ -597,6 +598,7 @@ impl Ingester {
let queue_id = subrequest.queue_id;

let batch_num_docs = subrequest.doc_batch.num_docs() as u64;
let batch_size = doc_batch_size(&subrequest.doc_batch, force_commit);

let append_result = append_non_empty_doc_batch(
&mut state_guard.mrecordlog,
Expand Down Expand Up @@ -637,11 +639,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 += batch_size;

let persist_success = PersistSuccess {
subrequest_id: subrequest.subrequest_id,
Expand Down Expand Up @@ -1199,6 +1202,7 @@ mod tests {
use crate::ingest_v2::DEFAULT_IDLE_SHARD_TIMEOUT;
use crate::ingest_v2::doc_mapper::try_build_doc_mapper;
use crate::ingest_v2::fetch::tests::{into_fetch_eof, into_fetch_payload};
use crate::ingest_v2::mrecordlog_utils::read_queue_size;

pub(super) struct IngesterForTest {
node_id: NodeId,
Expand Down Expand Up @@ -1782,6 +1786,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,
read_queue_size(&state_guard.mrecordlog, &queue_id_01)
);

state_guard.mrecordlog.assert_records_eq(
&queue_id_01,
Expand All @@ -1793,6 +1801,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,
read_queue_size(&state_guard.mrecordlog, &queue_id_11)
);

state_guard.mrecordlog.assert_records_eq(
&queue_id_11,
Expand Down
10 changes: 10 additions & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/models.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
use std::sync::Arc;
use std::time::{Duration, Instant};

use bytesize::ByteSize;
use quickwit_common::rate_limiter::RateLimiter;
use quickwit_doc_mapper::DocMapper;
use quickwit_proto::ingest::ShardState;
Expand All @@ -37,6 +38,7 @@ pub(super) struct IngesterShard {
pub replication_position_inclusive: Position,
/// Position up to which the shard has been truncated.
pub truncation_position_inclusive: Position,
pub queue_size: ByteSize,
pub rate_limiter: RateLimiter,
pub rate_meter: RateMeter,
/// Whether the shard should be advertised to other nodes (routers) via gossip.
Expand Down Expand Up @@ -67,6 +69,7 @@ pub(super) struct IngesterShardBuilder {
shard_state: ShardState,
replication_position_inclusive: Position,
truncation_position_inclusive: Position,
queue_size: ByteSize,
rate_limiter: RateLimiter,
rate_meter: RateMeter,
doc_mapper_opt: Option<Arc<DocMapper>>,
Expand Down Expand Up @@ -112,6 +115,11 @@ impl IngesterShardBuilder {
self
}

pub fn with_queue_size(mut self, queue_size: ByteSize) -> Self {
self.queue_size = queue_size;
self
}

/// Sets whether to validate documents. Defaults to `false`.
pub fn with_validate_docs(mut self, validate_docs: bool) -> Self {
self.validate_docs = validate_docs;
Expand Down Expand Up @@ -144,6 +152,7 @@ impl IngesterShardBuilder {
shard_state: self.shard_state,
replication_position_inclusive: self.replication_position_inclusive,
truncation_position_inclusive: self.truncation_position_inclusive,
queue_size: self.queue_size,
rate_limiter: self.rate_limiter,
rate_meter: self.rate_meter,
is_advertisable: self.is_advertisable,
Expand All @@ -170,6 +179,7 @@ impl IngesterShard {
shard_state: ShardState::Open,
replication_position_inclusive: Position::Beginning,
truncation_position_inclusive: Position::Beginning,
queue_size: ByteSize::default(),
rate_limiter: RateLimiter::default(),
rate_meter: RateMeter::default(),
doc_mapper_opt: None,
Expand Down
46 changes: 45 additions & 1 deletion quickwit/quickwit-ingest/src/ingest_v2/mrecordlog_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,9 @@ use quickwit_proto::ingest::DocBatchV2;
use quickwit_proto::types::{Position, QueueId};
use tracing::instrument;

use crate::MRecord;
use super::mrecord::MRECORD_HEADER_LEN;
use crate::mrecordlog_async::MultiRecordLogAsync;
use crate::{MRecord, estimate_size};

#[derive(Debug, thiserror::Error)]
pub(super) enum AppendDocBatchError {
Expand Down Expand Up @@ -98,6 +99,24 @@ pub(super) async fn append_non_empty_doc_batch(
}
}

pub(super) fn doc_batch_size(doc_batch: &DocBatchV2, force_commit: bool) -> ByteSize {
let estimated_size = estimate_size(doc_batch);
if force_commit {
estimated_size + ByteSize::b(MRECORD_HEADER_LEN as u64)
} else {
estimated_size
}
}

pub(super) fn read_queue_size(mrecordlog: &MultiRecordLogAsync, queue_id: &QueueId) -> ByteSize {
let num_bytes: usize = mrecordlog
.range(queue_id, ..)
.expect("queue should exist")
.map(|record| record.payload.len())
.sum();
ByteSize::b(num_bytes as u64)
}

/// Error returned when the mrecordlog does not have enough capacity to store some records.
#[derive(Debug, Clone, Copy, thiserror::Error)]
pub(super) enum NotEnoughCapacityError {
Expand Down Expand Up @@ -210,6 +229,31 @@ mod tests {
assert_eq!(position, Position::offset(2u64));
}

#[cfg(not(feature = "failpoints"))]
#[tokio::test]
async fn test_doc_batch_size_and_read_queue_size() {
let tempdir = tempfile::tempdir().unwrap();
let mut mrecordlog = MultiRecordLogAsync::open(tempdir.path()).await.unwrap();

let queue_id = "test-queue".to_string();
mrecordlog.create_queue(&queue_id).await.unwrap();
assert_eq!(read_queue_size(&mrecordlog, &queue_id), ByteSize::b(0));

let doc_batch = DocBatchV2::for_test(["test-doc-foo", "test-doc-bar"]);
assert_eq!(doc_batch_size(&doc_batch, false), ByteSize::b(28));
assert_eq!(doc_batch_size(&doc_batch, true), ByteSize::b(30));

append_non_empty_doc_batch(&mut mrecordlog, &queue_id, doc_batch.clone(), false)
.await
.unwrap();
assert_eq!(read_queue_size(&mrecordlog, &queue_id), ByteSize::b(28));

append_non_empty_doc_batch(&mut mrecordlog, &queue_id, doc_batch, true)
.await
.unwrap();
assert_eq!(read_queue_size(&mrecordlog, &queue_id), ByteSize::b(58));
}

// This test should be run manually and independently of other tests with the `failpoints`
// feature enabled:
// ```sh
Expand Down
16 changes: 15 additions & 1 deletion quickwit/quickwit-ingest/src/ingest_v2/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ use tracing::{error, info, instrument};
use super::local_shards::{ShardInfo, ShardInfos, ShardThroughputReadings};
use super::metrics::report_local_shards_metrics;
use super::models::IngesterShard;
use super::mrecordlog_utils::read_queue_size;
use super::rate_meter::RateMeter;
use super::wal_capacity_tracker::WalCapacityTracker;
use crate::OpenShardCounts;
Expand Down Expand Up @@ -392,6 +393,7 @@ impl IngesterState {
.checked_sub(1)
.map(Position::offset)
.unwrap_or(Position::Beginning);
let queue_size = read_queue_size(&mrecordlog, &queue_id);
let rate_limiter = RateLimiter::from_settings(rate_limiter_settings);
let rate_meter = RateMeter::default();

Expand All @@ -400,6 +402,7 @@ impl IngesterState {
.with_state(ShardState::Closed)
.with_replication_position_inclusive(replication_position_inclusive)
.with_truncation_position_inclusive(truncation_position_inclusive)
.with_queue_size(queue_size)
.with_rate_limiter(rate_limiter)
.with_rate_meter(rate_meter)
.with_last_write(now)
Expand Down Expand Up @@ -714,7 +717,13 @@ impl FullyLockedIngesterState<'_> {
.truncate(queue_id, truncate_up_to_offset_inclusive)
.await
{
Ok(_) => {}
Ok(evicted_size) => {
let queue_size = shard
.queue_size
.as_u64()
.saturating_sub(evicted_size.as_u64());
shard.queue_size = ByteSize::b(queue_size);
}
Err(TruncateError::MissingQueue(_)) => {
error!("failed to truncate shard `{queue_id}`: WAL queue not found");
self.shards.remove(queue_id);
Expand Down Expand Up @@ -1053,6 +1062,7 @@ mod tests {
shard_01.truncation_position_inclusive,
Position::offset(0u64)
);
assert_eq!(shard_01.queue_size, ByteSize::b(24));

// Fully truncated queue: recovers at its last position rather than the beginning.
let shard_02 = state_guard.shards.get(&queue_id_02).unwrap();
Expand All @@ -1065,12 +1075,14 @@ mod tests {
shard_02.truncation_position_inclusive,
Position::offset(1u64)
);
assert_eq!(shard_02.queue_size, ByteSize::b(0));

// Never-written queue: recovers at the beginning.
let shard_03 = state_guard.shards.get(&queue_id_03).unwrap();
assert_eq!(shard_03.shard_state, ShardState::Closed);
assert_eq!(shard_03.replication_position_inclusive, Position::Beginning);
assert_eq!(shard_03.truncation_position_inclusive, Position::Beginning);
assert_eq!(shard_03.queue_size, ByteSize::b(0));
}

fn insert_shard_with_used_capacity(
Expand Down Expand Up @@ -1361,6 +1373,7 @@ mod tests {
IngesterShard::builder(index_uid.clone(), source_id.clone(), ShardId::from(2))
.with_state(ShardState::Closed)
.with_replication_position_inclusive(Position::offset(1u64))
.with_queue_size(ByteSize::b(24))
.build();
state_guard.shards.insert(queue_id_02.clone(), shard_02);

Expand All @@ -1387,6 +1400,7 @@ mod tests {
.await;
let shard_02 = state_guard.shards.get(&queue_id_02).unwrap();
assert_eq!(shard_02.truncation_position_inclusive, Position::eof(1u64));
assert_eq!(shard_02.queue_size, ByteSize::b(0));
state_guard
.mrecordlog
.assert_records_eq(&queue_id_02, .., &[]);
Expand Down
9 changes: 7 additions & 2 deletions quickwit/quickwit-ingest/src/mrecordlog_async.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ use std::ops::RangeBounds;
use std::path::Path;

use bytes::Buf;
use bytesize::ByteSize;
use mrecordlog::error::*;
use mrecordlog::{MultiRecordLog, PersistAction, PersistPolicy, Record, ResourceUsage};
use tokio::task::JoinError;
Expand Down Expand Up @@ -151,7 +152,11 @@ impl MultiRecordLogAsync {
}

#[instrument(name = "mrecordlog.truncate_async", skip_all, fields(queue, position))]
pub async fn truncate(&mut self, queue: &str, position: u64) -> Result<usize, TruncateError> {
pub async fn truncate(
&mut self,
queue: &str,
position: u64,
) -> Result<ByteSize, TruncateError> {
let span = info_span!("mrecordlog.truncate", queue, position);
let queue = queue.to_string();
self.run_operation(span, move |mrecordlog| {
Expand All @@ -160,7 +165,7 @@ impl MultiRecordLogAsync {
.inspect(|outcome| {
WAL_BYTES_WRITTEN_TRUNCATE.inc_by(outcome.wal_bytes_written);
})
.map(|outcome| outcome.evicted_records)
.map(|outcome| ByteSize::b(outcome.evicted_bytes as u64))
})
.await
}
Expand Down
Loading