Skip to content
Merged
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/quickwit-cluster/src/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ impl Cluster {
];

if let Some(az) = &self_node.availability_zone {
initial_key_values.push((AVAILABILITY_ZONE_KEY.to_string(), az.clone()));
initial_key_values.push((AVAILABILITY_ZONE_KEY.to_string(), az.to_string()));
}
initial_key_values.push((
STANDALONE_COMPACTORS_KEY.to_string(),
Expand Down
10 changes: 5 additions & 5 deletions quickwit/quickwit-cluster/src/member.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use chitchat::{ChitchatId, NodeState, Version};
use quickwit_common::shared_consts::INGESTER_STATUS_KEY;
use quickwit_proto::indexing::{CpuCapacity, IndexingTask};
use quickwit_proto::ingest::ingester::IngesterStatus;
use quickwit_proto::types::NodeId;
use quickwit_proto::types::{AvailabilityZone, NodeId};
use tracing::{error, warn};

use crate::cluster::parse_indexing_tasks;
Expand Down Expand Up @@ -53,7 +53,7 @@ pub(crate) trait NodeStateExt {

fn ingester_status(&self) -> IngesterStatus;

fn availability_zone(&self) -> Option<String>;
fn availability_zone(&self) -> Option<AvailabilityZone>;

fn enable_standalone_compactors(&self) -> bool;
}
Expand Down Expand Up @@ -93,8 +93,8 @@ impl NodeStateExt for NodeState {
.unwrap_or(IngesterStatus::Ready)
}

fn availability_zone(&self) -> Option<String> {
self.get(AVAILABILITY_ZONE_KEY).map(|az| az.to_string())
fn availability_zone(&self) -> Option<AvailabilityZone> {
self.get(AVAILABILITY_ZONE_KEY).map(AvailabilityZone::from)
}

fn enable_standalone_compactors(&self) -> bool {
Expand Down Expand Up @@ -132,7 +132,7 @@ pub struct ClusterMember {
/// Whether the node is ready to serve requests.
pub is_ready: bool,
/// Availability zone the node is running in, if enabled.
pub availability_zone: Option<String>,
pub availability_zone: Option<AvailabilityZone>,
/// Whether the node was started with standalone compactors enabled.
pub enable_standalone_compactors: bool,
}
Expand Down
5 changes: 3 additions & 2 deletions quickwit/quickwit-cluster/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use quickwit_config::service::QuickwitService;
use quickwit_proto::indexing::IndexingTask;
#[cfg(any(test, feature = "testsuite"))]
use quickwit_proto::ingest::ingester::IngesterStatus;
use quickwit_proto::types::AvailabilityZone;
use tonic::transport::Channel;

use crate::member::{ClusterMember, build_cluster_member};
Expand Down Expand Up @@ -106,8 +107,8 @@ impl ClusterNode {
self.inner.is_self_node
}

pub fn availability_zone(&self) -> Option<&str> {
self.inner.member.availability_zone.as_deref()
pub fn availability_zone(&self) -> Option<AvailabilityZone> {
self.inner.member.availability_zone.clone()
}

pub fn enable_standalone_compactors(&self) -> bool {
Expand Down
4 changes: 2 additions & 2 deletions quickwit/quickwit-config/src/node_config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ use quickwit_common::shared_consts::{
use quickwit_common::uri::Uri;
use quickwit_proto::indexing::CpuCapacity;
use quickwit_proto::tonic::codec::CompressionEncoding;
use quickwit_proto::types::NodeId;
use quickwit_proto::types::{AvailabilityZone, NodeId};
use serde::{Deserialize, Deserializer, Serialize};
use tracing::{info, warn};

Expand Down Expand Up @@ -916,7 +916,7 @@ impl Default for JaegerConfig {
pub struct NodeConfig {
pub cluster_id: String,
pub node_id: NodeId,
pub availability_zone: Option<String>,
pub availability_zone: Option<AvailabilityZone>,
pub enabled_services: HashSet<QuickwitService>,
pub gossip_listen_addr: SocketAddr,
pub grpc_listen_addr: SocketAddr,
Expand Down
22 changes: 14 additions & 8 deletions quickwit/quickwit-config/src/node_config/serialize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ use quickwit_common::fs::get_disk_size;
use quickwit_common::net::{Host, find_private_ip, get_short_hostname};
use quickwit_common::new_coolid;
use quickwit_common::uri::Uri;
use quickwit_proto::types::NodeId;
use quickwit_proto::types::{AvailabilityZone, NodeId};
use serde::{Deserialize, Serialize};
use tracing::{info, warn};

Expand Down Expand Up @@ -255,11 +255,17 @@ impl NodeConfigBuilder {
.node_id
.resolve(env_vars)
.map(|node_id_str| NodeId::from_str(&node_id_str))?;
let availability_zone = self
.availability_zone
.resolve_optional(env_vars)?
.map(|availability_zone| availability_zone.trim().to_string())
.filter(|availability_zone| !availability_zone.is_empty());
let availability_zone =
self.availability_zone
.resolve_optional(env_vars)?
.and_then(|availability_zone| {
let availability_zone = availability_zone.trim();
if availability_zone.is_empty() {
None
} else {
Some(AvailabilityZone::from(availability_zone))
}
});

let enable_standalone_compactors = self.enable_standalone_compactors.resolve(env_vars)?;
let docs_clustering_config =
Expand Down Expand Up @@ -633,7 +639,7 @@ pub fn node_config_for_tests_from_ports(
) -> NodeConfig {
let node_id = NodeId::from_str(&default_node_id().unwrap());
let enabled_services = QuickwitService::default_services();
let availability_zone = Some(String::from("az-1"));
let availability_zone = Some(AvailabilityZone::from("az-1"));
let listen_address = Host::default();
let rest_listen_addr = listen_address
.with_port(rest_listen_port)
Expand Down Expand Up @@ -729,7 +735,7 @@ mod tests {
assert!(config.is_service_enabled(QuickwitService::Janitor));
assert!(config.is_service_enabled(QuickwitService::Metastore));

assert_eq!(config.availability_zone.unwrap(), "az-1");
assert_eq!(config.availability_zone.as_deref(), Some("az-1"));
assert_eq!(
config.rest_config.listen_addr,
SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 1111)
Expand Down
17 changes: 7 additions & 10 deletions quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ use quickwit_proto::indexing::{
use quickwit_proto::ingest::ingester::IngesterStatus;
use quickwit_proto::types::NodeId;
use scheduling::{
AvailabilityZone, Eligibility, IndexerInfo, SourceToSchedule, SourceToScheduleType,
Eligibility, IndexerInfo, SourceToSchedule, SourceToScheduleType,
compute_max_num_shards_per_pipeline, is_shard_in_same_zone,
};
use serde::Serialize;
Expand Down Expand Up @@ -332,10 +332,7 @@ fn build_indexer_info(indexer: &IndexerPoolEntry, locality_aware: bool) -> Index
};
IndexerInfo {
cpu_capacity: indexer.indexing_capacity,
availability_zone: indexer
.availability_zone
.as_deref()
.map(AvailabilityZone::from),
availability_zone: indexer.availability_zone.clone(),
eligibility,
}
}
Expand Down Expand Up @@ -983,7 +980,7 @@ mod tests {
use proptest::{prop_compose, proptest};
use quickwit_config::{IndexConfig, KafkaSourceParams, SourceConfig, SourceParams};
use quickwit_metastore::IndexMetadata;
use quickwit_proto::types::{IndexUid, PipelineUid, ShardId, SourceUid};
use quickwit_proto::types::{AvailabilityZone, IndexUid, PipelineUid, ShardId, SourceUid};

use super::*;
use crate::indexing_scheduler::scheduling::{
Expand Down Expand Up @@ -1685,7 +1682,7 @@ mod tests {
fn test_all_indexers_advertise_availability_zone() {
let indexer_pool = IndexerPool::default();
let mut zoned_indexer = mock_indexer_node_info("indexer-zoned", IngesterStatus::Ready);
zoned_indexer.availability_zone = Some("az-a".to_string());
zoned_indexer.availability_zone = Some(AvailabilityZone::from("az-a"));
indexer_pool.insert(zoned_indexer.node_id.clone(), zoned_indexer);

assert!(all_indexers_advertise_availability_zone(&indexer_pool));
Expand All @@ -1702,7 +1699,7 @@ mod tests {
let locality_aware = true;
{
let mut ready = mock_indexer_node_info("indexer-ready", IngesterStatus::Ready);
ready.availability_zone = Some("az-a".to_string());
ready.availability_zone = Some(AvailabilityZone::from("az-a"));
let retiring = mock_indexer_node_info("indexer-retiring", IngesterStatus::Retiring);
let decommissioning =
mock_indexer_node_info("indexer-decommissioning", IngesterStatus::Decommissioning);
Expand Down Expand Up @@ -1743,9 +1740,9 @@ mod tests {
}
{
let mut ready = mock_indexer_node_info("indexer-ready", IngesterStatus::Ready);
ready.availability_zone = Some("az-a".to_string());
ready.availability_zone = Some(AvailabilityZone::from("az-a"));
let mut retiring = mock_indexer_node_info("indexer-retiring", IngesterStatus::Retiring);
retiring.availability_zone = Some("az-b".to_string());
retiring.availability_zone = Some(AvailabilityZone::from("az-b"));
let indexers = vec![ready, retiring];
let locality_unaware = false;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,11 @@ pub mod scheduling_logic_model;

use std::collections::HashMap;
use std::num::NonZeroU32;
use std::sync::Arc;

use fnv::{FnvHashMap, FnvHashSet};
use quickwit_common::rate_limited_debug;
use quickwit_proto::indexing::{CpuCapacity, IndexingTask};
use quickwit_proto::types::{NodeId, PipelineUid, ShardId, SourceUid};
use quickwit_proto::types::{AvailabilityZone, NodeId, PipelineUid, ShardId, SourceUid};
pub use scheduling_logic_model::Eligibility;
use scheduling_logic_model::{IndexerLocality, IndexerOrd, LocalityGroup, SourceOrd};
use tracing::{error, warn};
Expand All @@ -38,8 +37,6 @@ use crate::indexing_scheduler::scheduling::scheduling_logic_model::{
};
use crate::model::ShardLocations;

pub type AvailabilityZone = Arc<str>;

/// If we have several pipelines below this threshold we
/// reduce the number of pipelines.
///
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,13 @@ use std::time::Instant;
use fnv::FnvHashMap;
use quickwit_proto::indexing::IndexingTask;
use quickwit_proto::ingest::ingester::IngesterStatus;
use quickwit_proto::types::{NodeId, ShardId, SourceUid};
use quickwit_proto::types::{AvailabilityZone, NodeId, ShardId, SourceUid};
use rand::Rng;
use rand::seq::SliceRandom;

use super::{
AvailabilityZone, Eligibility, IndexerInfo, SourceToSchedule,
compute_max_num_shards_per_pipeline, shard_availability_zone,
Eligibility, IndexerInfo, SourceToSchedule, compute_max_num_shards_per_pipeline,
shard_availability_zone,
};
use crate::IndexerPoolEntry;
use crate::indexing_plan::PhysicalIndexingPlan;
Expand Down Expand Up @@ -364,7 +364,9 @@ mod tests {
use fnv::FnvHashMap;
use quickwit_proto::indexing::{IndexingTask, mcpu};
use quickwit_proto::ingest::ingester::IngesterStatus;
use quickwit_proto::types::{IndexUid, NodeId, PipelineUid, ShardId, SourceUid};
use quickwit_proto::types::{
AvailabilityZone, IndexUid, NodeId, PipelineUid, ShardId, SourceUid,
};
use rand::SeedableRng;
use rand::rngs::StdRng;

Expand All @@ -376,8 +378,7 @@ mod tests {
};
use crate::indexing_plan::PhysicalIndexingPlan;
use crate::indexing_scheduler::scheduling::{
AvailabilityZone, Eligibility, IndexerInfo, SourceToSchedule, SourceToScheduleType,
shard_ids_for_indexer,
Eligibility, IndexerInfo, SourceToSchedule, SourceToScheduleType, shard_ids_for_indexer,
};
use crate::indexing_scheduler::{
IndexingSchedulerState, MIN_DURATION_BETWEEN_SCHEDULING, get_indexing_plan_density,
Expand Down
Loading
Loading