From b1127946b2285f5caea20c97f450d33d9c71446d Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Thu, 8 Oct 2026 13:22:52 -0400 Subject: [PATCH] Refactor BroadcastLocalShards task to use new shard readings --- .../src/ingest_v2/broadcast/local_shards.rs | 297 +++++------------- .../quickwit-ingest/src/ingest_v2/ingester.rs | 27 +- .../quickwit-ingest/src/ingest_v2/models.rs | 17 +- .../src/ingest_v2/rate_meter.rs | 14 +- .../src/ingest_v2/shard_readings.rs | 62 +++- .../quickwit-ingest/src/ingest_v2/state.rs | 4 +- 6 files changed, 164 insertions(+), 257 deletions(-) diff --git a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs index 4d22f8a0fa6..e298e230d3f 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/broadcast/local_shards.rs @@ -12,27 +12,24 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::{BTreeMap, BTreeSet, HashMap}; -use std::time::Duration; +use std::collections::{BTreeMap, BTreeSet}; +use std::sync::Arc; use bytesize::ByteSize; use quickwit_cluster::{Cluster, ListenerHandle}; use quickwit_common::pubsub::{Event, EventBroker}; -use quickwit_common::ring_buffer::RingBuffer; use quickwit_common::shared_consts::INGESTER_SHARDS_PREFIX; use quickwit_common::sorted_iter::{KeyDiff, SortedByKeyIterator}; -use quickwit_common::tower::{ConstantRate, Rate}; use quickwit_proto::ingest::ShardState; use quickwit_proto::types::{NodeId, ShardId, SourceUid}; use serde::{Deserialize, Serialize, Serializer}; +use tokio::sync::watch; use tokio::task::JoinHandle; use tracing::{debug, instrument, warn}; use super::{BROADCAST_INTERVAL_PERIOD, make_key, parse_key}; use crate::RateMibPerSec; -use crate::ingest_v2::metrics::{ - CLOSED_SHARDS, OPEN_SHARDS, SHARD_LT_THROUGHPUT_MIB, SHARD_ST_THROUGHPUT_MIB, -}; +use crate::ingest_v2::shard_readings::{ShardThroughputReading, ShardThroughputReadings}; use crate::ingest_v2::state::WeakIngesterState; const ONE_MIB: ByteSize = ByteSize::mib(1); @@ -101,6 +98,29 @@ impl<'de> Deserialize<'de> for ShardInfo { } } +impl From<&ShardThroughputReading> for ShardInfo { + fn from(reading: &ShardThroughputReading) -> Self { + let short_term_ingestion_rate_mib_per_sec_u64: u64 = reading + .short_term_ingestion_rate + .as_u64() + .div_ceil(ONE_MIB.as_u64()); + let long_term_ingestion_rate_mib_per_sec_u64: u64 = reading + .long_term_ingestion_rate + .as_u64() + .div_ceil(ONE_MIB.as_u64()); + Self { + shard_id: reading.shard_id.clone(), + shard_state: reading.shard_state, + short_term_ingestion_rate: RateMibPerSec( + short_term_ingestion_rate_mib_per_sec_u64 as u16, + ), + long_term_ingestion_rate: RateMibPerSec( + long_term_ingestion_rate_mib_per_sec_u64 as u16, + ), + } + } +} + /// A set of shards belonging to the same source. pub type ShardInfos = BTreeSet; @@ -121,6 +141,19 @@ enum ShardInfosChange<'a> { }, } +impl From<&ShardThroughputReadings> for LocalShardsSnapshot { + fn from(readings: &ShardThroughputReadings) -> Self { + let mut per_source_shard_infos = BTreeMap::new(); + for (source_uid, shard_readings) in &readings.per_source_readings { + let shard_infos: ShardInfos = shard_readings.iter().map(ShardInfo::from).collect(); + per_source_shard_infos.insert(source_uid.clone(), shard_infos); + } + Self { + per_source_shard_infos, + } + } +} + impl LocalShardsSnapshot { pub fn diff<'a>(&'a self, other: &'a Self) -> impl Iterator> + 'a { self.per_source_shard_infos @@ -153,164 +186,27 @@ impl LocalShardsSnapshot { pub struct BroadcastLocalShardsTask { cluster: Cluster, weak_state: WeakIngesterState, - shard_throughput_time_series_map: ShardThroughputTimeSeriesMap, + local_shards_rx: watch::Receiver>>, /// Snapshot broadcast on the previous tick. Carried across iterations so /// we can diff against the new snapshot and only broadcast changes. previous_snapshot: LocalShardsSnapshot, } -const SHARD_THROUGHPUT_LONG_TERM_WINDOW_LEN: usize = 12; - -#[derive(Default)] -struct ShardThroughputTimeSeriesMap { - shard_time_series: HashMap<(SourceUid, ShardId), ShardThroughputTimeSeries>, -} - -impl ShardThroughputTimeSeriesMap { - // Records a list of shard throughputs. - // - // A new time series is created for each new shard_ids. - // If a shard_id had a time series, and it is not present in the - // `shard_throughput`, the time series will be removed. - #[allow(clippy::mutable_key_type)] - pub fn record_shard_throughputs( - &mut self, - shard_throughputs: HashMap<(SourceUid, ShardId), (ShardState, ConstantRate)>, - ) { - self.shard_time_series - .retain(|key, _| shard_throughputs.contains_key(key)); - for ((source_uid, shard_id), (shard_state, throughput)) in shard_throughputs { - let throughput_measurement = throughput.rescale(Duration::from_secs(1)).work_bytes(); - let shard_time_series = self - .shard_time_series - .entry((source_uid.clone(), shard_id.clone())) - .or_default(); - shard_time_series.shard_state = shard_state; - shard_time_series.record(throughput_measurement); - } - } - - pub fn get_per_source_shard_infos(&self) -> BTreeMap { - let mut per_source_shard_infos: BTreeMap = BTreeMap::new(); - for ((source_uid, shard_id), shard_time_series) in self.shard_time_series.iter() { - let shard_state = shard_time_series.shard_state; - let short_term_ingestion_rate_mib_per_sec_u64: u64 = - shard_time_series.last().as_u64().div_ceil(ONE_MIB.as_u64()); - let long_term_ingestion_rate_mib_per_sec_u64: u64 = shard_time_series - .average() - .as_u64() - .div_ceil(ONE_MIB.as_u64()); - SHARD_ST_THROUGHPUT_MIB.observe(short_term_ingestion_rate_mib_per_sec_u64 as f64); - SHARD_LT_THROUGHPUT_MIB.observe(long_term_ingestion_rate_mib_per_sec_u64 as f64); - - let short_term_ingestion_rate = - RateMibPerSec(short_term_ingestion_rate_mib_per_sec_u64 as u16); - let long_term_ingestion_rate = - RateMibPerSec(long_term_ingestion_rate_mib_per_sec_u64 as u16); - let shard_info = ShardInfo { - shard_id: shard_id.clone(), - shard_state, - short_term_ingestion_rate, - long_term_ingestion_rate, - }; - - per_source_shard_infos - .entry(source_uid.clone()) - .or_default() - .insert(shard_info); - } - per_source_shard_infos - } -} - -#[derive(Default)] -struct ShardThroughputTimeSeries { - shard_state: ShardState, - throughput: RingBuffer, -} - -impl ShardThroughputTimeSeries { - fn last(&self) -> ByteSize { - self.throughput.last().unwrap_or_default() - } - - fn average(&self) -> ByteSize { - if self.throughput.is_empty() { - return ByteSize::default(); - } - let sum = self.throughput.iter().map(ByteSize::as_u64).sum::(); - ByteSize::b(sum / self.throughput.len() as u64) - } - - fn record(&mut self, new_throughput_measurement: ByteSize) { - self.throughput.push_back(new_throughput_measurement); - } -} - impl BroadcastLocalShardsTask { - pub fn spawn(cluster: Cluster, weak_state: WeakIngesterState) -> JoinHandle<()> { + pub fn spawn( + cluster: Cluster, + weak_state: WeakIngesterState, + local_shards_rx: watch::Receiver>>, + ) -> JoinHandle<()> { let broadcaster = Self { cluster, weak_state, - shard_throughput_time_series_map: Default::default(), + local_shards_rx, previous_snapshot: LocalShardsSnapshot::default(), }; tokio::spawn(broadcaster.run()) } - async fn snapshot_local_shards(&mut self) -> Option { - let state = self.weak_state.upgrade()?; - - let Ok(mut state_guard) = state.lock_partially("snapshot_local_shards").await else { - return Some(LocalShardsSnapshot::default()); - }; - #[allow(clippy::mutable_key_type)] - let ingestion_rates: HashMap<(SourceUid, ShardId), (ShardState, ConstantRate)> = - state_guard - .shards - .values_mut() - .filter(|shard| shard.is_advertisable) - .map(|shard| { - let source_uid = SourceUid { - index_uid: shard.index_uid.clone(), - source_id: shard.source_id.clone(), - }; - let shard_id = shard.shard_id.clone(); - let shard_state = shard.shard_state; - let rate_meter = &mut shard.rate_meter; - - ((source_uid, shard_id), (shard_state, rate_meter.harvest())) - }) - .collect(); - - self.shard_throughput_time_series_map - .record_shard_throughputs(ingestion_rates); - - let per_source_shard_infos = self - .shard_throughput_time_series_map - .get_per_source_shard_infos(); - - let mut num_open_shards = 0; - let mut num_closed_shards = 0; - - for shard_infos in per_source_shard_infos.values() { - for shard_info in shard_infos { - match shard_info.shard_state { - ShardState::Open => num_open_shards += 1, - ShardState::Closed => num_closed_shards += 1, - ShardState::Unavailable | ShardState::Unspecified => {} - } - } - } - OPEN_SHARDS.set(num_open_shards as f64); - CLOSED_SHARDS.set(num_closed_shards as f64); - - let snapshot = LocalShardsSnapshot { - per_source_shard_infos, - }; - Some(snapshot) - } - async fn broadcast_local_shards( &self, previous_snapshot: &LocalShardsSnapshot, @@ -353,9 +249,13 @@ impl BroadcastLocalShardsTask { /// should stop. #[instrument(name = "broadcast_local_shards.tick", skip_all)] async fn run_once(&mut self) -> bool { - let Some(new_snapshot) = self.snapshot_local_shards().await else { + if self.weak_state.upgrade().is_none() { return false; + } + let Some(readings) = self.local_shards_rx.borrow().clone() else { + return true; }; + let new_snapshot = LocalShardsSnapshot::from(readings.as_ref()); self.broadcast_local_shards(&self.previous_snapshot, &new_snapshot) .await; self.previous_snapshot = new_snapshot; @@ -401,8 +301,8 @@ pub async fn setup_local_shards_update_listener( #[cfg(test)] mod tests { - use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::Duration; use quickwit_cluster::{ChitchatTransport, create_cluster_for_test}; use quickwit_common::shared_consts::INGESTER_SHARDS_PREFIX; @@ -411,7 +311,6 @@ mod tests { use super::*; use crate::RateMibPerSec; - use crate::ingest_v2::models::IngesterShard; use crate::ingest_v2::state::IngesterState; #[test] @@ -541,59 +440,46 @@ mod tests { .await .unwrap(); let (_temp_dir, state) = IngesterState::for_test(cluster.clone()).await; + let (local_shards_tx, local_shards_rx) = watch::channel(None); let mut task = BroadcastLocalShardsTask { cluster, weak_state: state.weak(), - shard_throughput_time_series_map: Default::default(), + local_shards_rx, previous_snapshot: LocalShardsSnapshot::default(), }; - let mut state_guard = state.lock_partially("test").await.unwrap(); + // No readings published yet: nothing to broadcast. + assert!(task.run_once().await); + assert!(task.previous_snapshot.per_source_shard_infos.is_empty()); let index_uid = IndexUid::for_test("test-index", 0); - let shard_00 = IngesterShard::builder( - index_uid.clone(), - SourceId::from("test-source"), - ShardId::from(0), - ) - .build(); - state_guard.shards.insert(shard_00.queue_id(), shard_00); - - let shard_01 = IngesterShard::builder( - index_uid.clone(), - SourceId::from("test-source"), - ShardId::from(1), - ) - .advertisable() - .build(); - let queue_id_01 = shard_01.queue_id(); - state_guard.shards.insert(queue_id_01.clone(), shard_01); - - drop(state_guard); - - // First tick: shard_01 (advertisable) is the only one contributing to the snapshot — - // broadcast publishes it. + let source_uid = SourceUid { + index_uid: index_uid.clone(), + source_id: SourceId::from("test-source"), + }; + let reading = ShardThroughputReading { + shard_id: ShardId::from(1), + shard_state: ShardState::Open, + short_term_ingestion_rate: ByteSize::b(1), + long_term_ingestion_rate: ByteSize::mib(2), + }; + let readings = ShardThroughputReadings { + per_source_readings: BTreeMap::from([(source_uid, vec![reading])]), + }; + local_shards_tx.send_replace(Some(Arc::new(readings))); assert!(task.run_once().await); - assert_eq!(task.previous_snapshot.per_source_shard_infos.len(), 1); - - tokio::time::sleep(Duration::from_millis(100)).await; let key = format!("{INGESTER_SHARDS_PREFIX}{}:{}", index_uid, "test-source"); - task.cluster.get_self_key_value(&key).await.unwrap(); - - // Remove the only advertisable shard, run again: snapshot empty, - // broadcast clears the chitchat key. - let mut state_guard = state.lock_partially("test").await.unwrap(); - state_guard.shards.remove(&queue_id_01); - drop(state_guard); + let value = task.cluster.get_self_key_value(&key).await.unwrap(); + // Rates are rounded up to the next MiB/s. + assert_eq!(value, r#"["00000000000000000001:open:1:2"]"#); + local_shards_tx.send_replace(Some(Arc::new(ShardThroughputReadings::default()))); assert!(task.run_once().await); - assert!(task.previous_snapshot.per_source_shard_infos.is_empty()); + assert!(task.cluster.get_self_key_value(&key).await.is_none()); - tokio::time::sleep(Duration::from_millis(100)).await; - - let value_opt = task.cluster.get_self_key_value(&key).await; - assert!(value_opt.is_none()); + drop(state); + assert!(!task.run_once().await); } #[tokio::test] @@ -646,29 +532,4 @@ mod tests { assert_eq!(local_shards_update_counter.load(Ordering::Acquire), 1); } - - #[test] - fn test_shard_throughput_time_series() { - let mut time_series = ShardThroughputTimeSeries::default(); - assert_eq!(time_series.last(), ByteSize::mb(0)); - assert_eq!(time_series.average(), ByteSize::mb(0)); - - time_series.record(ByteSize::mb(2)); - assert_eq!(time_series.last(), ByteSize::mb(2)); - assert_eq!(time_series.average(), ByteSize::mb(2)); - - time_series.record(ByteSize::mb(1)); - assert_eq!(time_series.last(), ByteSize::mb(1)); - assert_eq!(time_series.average(), ByteSize::kb(1500)); - - time_series.record(ByteSize::mb(3)); - assert_eq!(time_series.last(), ByteSize::mb(3)); - assert_eq!(time_series.average(), ByteSize::mb(2)); - - for _ in 0..SHARD_THROUGHPUT_LONG_TERM_WINDOW_LEN { - time_series.record(ByteSize::mb(4)); - assert_eq!(time_series.last(), ByteSize::mb(4)); - } - assert_eq!(time_series.last(), ByteSize::mb(4)); - } } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs index 35dc2db2fc5..9c3cdb7ae7c 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/ingester.rs @@ -53,7 +53,6 @@ use super::models::IngesterShard; use super::mrecordlog_utils::{ AppendDocBatchError, check_enough_capacity, doc_batch_size, wal_stats, }; -use super::rate_meter::RateMeter; use super::shard_readings::{ShardReadingsPublisher, ShardThroughputReadings}; use super::state::{IngesterState, InnerIngesterState, WeakIngesterState}; use crate::ingest_v2::doc_mapper::get_or_try_build_doc_mapper; @@ -183,7 +182,11 @@ impl Ingester { .await; let weak_state = state.weak(); - BroadcastLocalShardsTask::spawn(cluster.clone(), weak_state.clone()); + BroadcastLocalShardsTask::spawn( + cluster.clone(), + weak_state.clone(), + local_shards_tx.subscribe(), + ); ShardReadingsPublisher::spawn(weak_state.clone(), local_shards_tx); BroadcastIngesterCapacityScoreTask::spawn(cluster, weak_state.clone()); CloseIdleShardsTask::spawn(weak_state, idle_shard_timeout); @@ -247,11 +250,9 @@ impl Ingester { let source_id = shard.source_id.clone(); let shard_id = shard.shard_id().clone(); let rate_limiter = RateLimiter::from_settings(self.rate_limiter_settings); - let rate_meter = RateMeter::default(); let shard = IngesterShard::builder(index_uid, source_id, shard_id) .with_rate_limiter(rate_limiter) - .with_rate_meter(rate_meter) .with_shared_rate_meter(state.shared_rate_meter.clone()) .with_doc_mapper(doc_mapper) .with_validate_docs(validate_docs) @@ -571,7 +572,6 @@ impl Ingester { ) .inc_by(original_batch_num_bytes - valid_batch_num_bytes); } - shard.rate_meter.update(valid_batch_num_bytes); total_requested_capacity += requested_capacity; let pending_persist_subrequest = PendingPersistSubrequest { @@ -1413,17 +1413,20 @@ mod tests { let source_id = SourceId::from("test-source"); let shard_00 = - IngesterShard::builder(index_uid.clone(), source_id.clone(), ShardId::from(0)).build(); + IngesterShard::builder(index_uid.clone(), source_id.clone(), ShardId::from(0)) + .with_shared_rate_meter(state_guard.shared_rate_meter.clone()) + .build(); state_guard.shards.insert(shard_00.queue_id(), shard_00); let shard_01 = IngesterShard::builder(index_uid.clone(), source_id, ShardId::from(1)) + .with_shared_rate_meter(state_guard.shared_rate_meter.clone()) .advertisable() .build(); let queue_id_01 = shard_01.queue_id(); state_guard.shards.insert(queue_id_01.clone(), shard_01); drop(state_guard); - tokio::time::sleep(Duration::from_millis(100)).await; + tokio::time::sleep(Duration::from_millis(200)).await; let key = format!("{INGESTER_SHARDS_PREFIX}{}:{}", index_uid, "test-source"); let value = ingester_ctx.cluster.get_self_key_value(&key).await.unwrap(); @@ -1437,14 +1440,10 @@ mod tests { assert_eq!(shard_info.short_term_ingestion_rate, 0); let mut state_guard = ingester.state.lock_fully("test").await.unwrap(); - state_guard - .shards - .get_mut(&queue_id_01) - .unwrap() - .shard_state = ShardState::Closed; + state_guard.shards.get_mut(&queue_id_01).unwrap().close(); drop(state_guard); - tokio::time::sleep(Duration::from_millis(100)).await; + tokio::time::sleep(Duration::from_millis(200)).await; let value = ingester_ctx.cluster.get_self_key_value(&key).await.unwrap(); @@ -1458,7 +1457,7 @@ mod tests { state_guard.shards.remove(&queue_id_01).unwrap(); drop(state_guard); - tokio::time::sleep(Duration::from_millis(100)).await; + tokio::time::sleep(Duration::from_millis(200)).await; let value_opt = ingester_ctx.cluster.get_self_key_value(&key).await; assert!(value_opt.is_none()); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/models.rs b/quickwit/quickwit-ingest/src/ingest_v2/models.rs index 1d1e685ac6a..a7b4c4e1985 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/models.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/models.rs @@ -23,7 +23,7 @@ use quickwit_proto::types::{IndexUid, Position, QueueId, ShardId, SourceId, Sour use tokio::sync::watch; use tracing::error; -use crate::ingest_v2::rate_meter::{RateMeter, SharedRateMeter}; +use crate::ingest_v2::rate_meter::SharedRateMeter; /// Status of a shard: state + position of the last record written. pub(super) type ShardStatus = (ShardState, Position); @@ -42,7 +42,6 @@ pub(super) struct IngesterShard { /// indexed. pub queue_size: ByteSize, pub rate_limiter: RateLimiter, - pub rate_meter: RateMeter, /// The shared rate meter contains throughput and status readings for all shards, centralized /// to be able to report to the control plane. shared_rate_meter: Arc, @@ -65,8 +64,9 @@ pub(super) struct IngesterShard { } /// Builder for `IngesterShard`. By default, the shard is open, is empty (i.e. the replication and -/// truncation positions are at the beginning), has a zero queue size, uses the default rate limiter -/// and rate meter, has no doc mapper, does not validate documents, and is not advertisable. +/// truncation positions are at the beginning), has a zero queue size and uses the default rate +/// limiter. A shared rate meter should be wired up to accurately capture readings. +/// It also has no doc mapper, does not validate documents, and is not advertisable. pub(super) struct IngesterShardBuilder { index_uid: IndexUid, source_id: SourceId, @@ -76,7 +76,6 @@ pub(super) struct IngesterShardBuilder { truncation_position_inclusive: Position, queue_size: ByteSize, rate_limiter: RateLimiter, - rate_meter: RateMeter, shared_rate_meter: Arc, doc_mapper_opt: Option>, validate_docs: bool, @@ -103,12 +102,6 @@ impl IngesterShardBuilder { self } - /// Sets the rate meter. Defaults to `RateMeter::default()`. - pub fn with_rate_meter(mut self, rate_meter: RateMeter) -> Self { - self.rate_meter = rate_meter; - self - } - pub fn with_shared_rate_meter(mut self, shared_rate_meter: Arc) -> Self { self.shared_rate_meter = shared_rate_meter; self @@ -176,7 +169,6 @@ impl IngesterShardBuilder { truncation_position_inclusive: self.truncation_position_inclusive, queue_size: self.queue_size, rate_limiter: self.rate_limiter, - rate_meter: self.rate_meter, shared_rate_meter: self.shared_rate_meter, is_advertisable: self.is_advertisable, doc_mapper_opt: self.doc_mapper_opt, @@ -204,7 +196,6 @@ impl IngesterShard { truncation_position_inclusive: Position::Beginning, queue_size: ByteSize::default(), rate_limiter: RateLimiter::default(), - rate_meter: RateMeter::default(), shared_rate_meter: Arc::new(SharedRateMeter::default()), doc_mapper_opt: None, validate_docs: false, diff --git a/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs b/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs index 649b966bc00..ce7777f5c7c 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/rate_meter.rs @@ -149,15 +149,15 @@ impl SharedRateMeter { } } -pub(super) struct IngestionRates { - pub short_term: ByteSize, - pub long_term: ByteSize, +struct IngestionRates { + short_term: ByteSize, + long_term: ByteSize, } /// A naive rate meter that tracks how much work was performed during a period of time defined by /// two successive calls to `harvest`. #[derive(Debug)] -pub(super) struct RateMeter { +struct RateMeter { total_work: u64, harvested_at: Instant, short_term_rates: RingBuffer, @@ -177,13 +177,13 @@ impl Default for RateMeter { impl RateMeter { /// Increments the amount of work performed since the last call to `harvest`. - pub fn update(&mut self, work: u64) { + fn update(&mut self, work: u64) { self.total_work += work; } /// Returns the average work rate since the last call to this method and resets the internal /// state. - pub fn harvest(&mut self) -> ConstantRate { + fn harvest(&mut self) -> ConstantRate { let now = Instant::now(); let elapsed = now.duration_since(self.harvested_at); let rate = ConstantRate::new(self.total_work, elapsed); @@ -192,7 +192,7 @@ impl RateMeter { rate } - pub fn sample(&mut self) -> IngestionRates { + fn sample(&mut self) -> IngestionRates { let rate = self.harvest(); let rate_per_sec = rate.rescale(Duration::from_secs(1)).work_bytes(); self.short_term_rates.push_back(rate_per_sec); diff --git a/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs b/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs index 280747a3bf2..663c36cde17 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/shard_readings.rs @@ -25,6 +25,9 @@ use tokio::task::JoinHandle; use tokio::time::MissedTickBehavior; use tracing::warn; +use super::metrics::{ + CLOSED_SHARDS, OPEN_SHARDS, SHARD_LT_THROUGHPUT_MIB, SHARD_ST_THROUGHPUT_MIB, +}; use super::state::WeakIngesterState; const SAMPLE_INTERVAL: Duration = if cfg!(any(test, feature = "testsuite")) { @@ -89,17 +92,41 @@ impl ShardReadingsPublisher { // possible, but unlikely if we're initialized but haven't sent a reading yet. continue; }; - let snapshot = Arc::new(shared_rate_meter.harvest()); - self.local_shards_tx.send_replace(Some(snapshot)); + let snapshot = shared_rate_meter.harvest(); + report_local_shards_metrics(&snapshot); + self.local_shards_tx.send_replace(Some(Arc::new(snapshot))); } } } +fn report_local_shards_metrics(snapshot: &ShardThroughputReadings) { + let mut num_open_shards = 0; + let mut num_closed_shards = 0; + + for shard_readings in snapshot.per_source_readings.values() { + for shard_reading in shard_readings { + match shard_reading.shard_state { + ShardState::Open => num_open_shards += 1, + ShardState::Closed => num_closed_shards += 1, + ShardState::Unavailable | ShardState::Unspecified => {} + } + SHARD_ST_THROUGHPUT_MIB + .observe(shard_reading.short_term_ingestion_rate.as_mib().ceil()); + SHARD_LT_THROUGHPUT_MIB.observe(shard_reading.long_term_ingestion_rate.as_mib().ceil()); + } + } + OPEN_SHARDS.set(num_open_shards as f64); + CLOSED_SHARDS.set(num_closed_shards as f64); +} + #[cfg(test)] mod tests { use quickwit_cluster::{ChitchatTransport, create_cluster_for_test}; + use quickwit_proto::types::IndexUid; use super::*; + use crate::ingest_v2::models::IngesterShard; + use crate::ingest_v2::rate_meter::SharedRateMeter; use crate::ingest_v2::state::IngesterState; async fn state() -> (tempfile::TempDir, IngesterState) { @@ -152,4 +179,35 @@ mod tests { tick().await; publisher.await.unwrap(); } + + #[test] + fn test_report_local_shards_metrics() { + let shared_rate_meter = Arc::new(SharedRateMeter::default()); + let mut shards = Vec::new(); + for (shard_id, shard_state) in [ + (1, ShardState::Open), + (2, ShardState::Open), + (3, ShardState::Closed), + ] { + let shard = IngesterShard::builder( + IndexUid::for_test("test-index", 0), + "test-source".to_string(), + ShardId::from(shard_id), + ) + .with_state(shard_state) + .with_shared_rate_meter(shared_rate_meter.clone()) + .advertisable() + .build(); + shards.push(shard); + } + report_local_shards_metrics(&shared_rate_meter.harvest()); + assert_eq!(OPEN_SHARDS.get(), 2.0); + assert_eq!(CLOSED_SHARDS.get(), 1.0); + + // Dropped shards are no longer counted. + drop(shards); + report_local_shards_metrics(&shared_rate_meter.harvest()); + assert_eq!(OPEN_SHARDS.get(), 0.0); + assert_eq!(CLOSED_SHARDS.get(), 0.0); + } } diff --git a/quickwit/quickwit-ingest/src/ingest_v2/state.rs b/quickwit/quickwit-ingest/src/ingest_v2/state.rs index afc1af727b6..ca27e1c0086 100644 --- a/quickwit/quickwit-ingest/src/ingest_v2/state.rs +++ b/quickwit/quickwit-ingest/src/ingest_v2/state.rs @@ -37,7 +37,7 @@ use tracing::{error, info, instrument}; use super::models::IngesterShard; use super::mrecordlog_utils::{AppendDocBatchError, append_non_empty_doc_batch}; -use super::rate_meter::{RateMeter, SharedRateMeter}; +use super::rate_meter::SharedRateMeter; use super::wal_capacity_tracker::WalCapacityTracker; use crate::OpenShardCounts; use crate::mrecordlog_async::MultiRecordLogAsync; @@ -309,7 +309,6 @@ impl IngesterState { .unwrap_or(Position::Beginning); let queue_size = ByteSize::b(queue_summary.num_bytes as u64); let rate_limiter = RateLimiter::from_settings(rate_limiter_settings); - let rate_meter = RateMeter::default(); let shard = IngesterShard::builder(index_uid.clone(), source_id.clone(), shard_id.clone()) @@ -318,7 +317,6 @@ impl IngesterState { .with_truncation_position_inclusive(truncation_position_inclusive) .with_queue_size(queue_size) .with_rate_limiter(rate_limiter) - .with_rate_meter(rate_meter) .with_shared_rate_meter(inner_guard.shared_rate_meter.clone()) .with_last_write(now) .advertisable() // We want to advertise the shard as read-only right away.