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
1 change: 1 addition & 0 deletions quickwit/quickwit-cli/src/tool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -965,6 +965,7 @@ async fn create_empty_cluster(config: &NodeConfig) -> anyhow::Result<Cluster> {
ingester_status: IngesterStatus::default(),
availability_zone: None,
enable_standalone_compactors: false,
enable_shard_scaling_v2: false,
};
let channel_factory = ChannelFactory::for_grpc(&config.grpc_config)?;
let cluster = Cluster::join(
Expand Down
7 changes: 6 additions & 1 deletion quickwit/quickwit-cluster/src/cluster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ use crate::grpc_gossip::spawn_catchup_callback_task;
use crate::member::{
AVAILABILITY_ZONE_KEY, ClusterMember, ENABLED_SERVICES_KEY, GRPC_ADVERTISE_ADDR_KEY,
NodeStateExt, PIPELINE_METRICS_PREFIX, READINESS_KEY, READINESS_VALUE_NOT_READY,
READINESS_VALUE_READY, STANDALONE_COMPACTORS_KEY,
READINESS_VALUE_READY, SHARD_SCALING_V2_KEY, STANDALONE_COMPACTORS_KEY,
};
use crate::metrics::spawn_metrics_task;
use crate::{ClusterChangeStream, ClusterNode};
Expand Down Expand Up @@ -241,6 +241,10 @@ impl Cluster {
STANDALONE_COMPACTORS_KEY.to_string(),
self_node.enable_standalone_compactors.to_string(),
));
initial_key_values.push((
SHARD_SCALING_V2_KEY.to_string(),
self_node.enable_shard_scaling_v2.to_string(),
));
let chitchat_handle =
spawn_chitchat(chitchat_config, initial_key_values, transport).await?;

Expand Down Expand Up @@ -776,6 +780,7 @@ impl<'a> TestClusterBuilder<'a> {
ingester_status: IngesterStatus::default(),
availability_zone: None,
enable_standalone_compactors: false,
enable_shard_scaling_v2: false,
},
peer_seed_addrs: Vec::new(),
transport,
Expand Down
7 changes: 5 additions & 2 deletions quickwit/quickwit-cluster/src/grpc_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ mod tests {
.key_values
.sort_unstable_by(|left, right| left.key.cmp(&right.key));

assert_eq!(node_state.key_values.len(), 5);
assert_eq!(node_state.key_values.len(), 6);
assert_eq!(node_state.key_values[0].key, ENABLED_SERVICES_KEY);
assert_eq!(node_state.key_values[0].value, "indexer");

Expand All @@ -179,8 +179,11 @@ mod tests {
assert_eq!(node_state.key_values[3].key, READINESS_KEY);
assert_eq!(node_state.key_values[3].value, "READY");

assert_eq!(node_state.key_values[4].key, STANDALONE_COMPACTORS_KEY);
assert_eq!(node_state.key_values[4].key, "shard_scaling_v2");
assert_eq!(node_state.key_values[4].value, "false");

assert_eq!(node_state.key_values[5].key, STANDALONE_COMPACTORS_KEY);
assert_eq!(node_state.key_values[5].value, "false");
}

#[tokio::test]
Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-cluster/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ pub async fn start_cluster_service(node_config: &NodeConfig) -> anyhow::Result<C
ingester_status: IngesterStatus::default(),
availability_zone: node_config.availability_zone.clone(),
enable_standalone_compactors: node_config.enable_standalone_compactors,
enable_shard_scaling_v2: node_config.enable_shard_scaling_v2,
};
let failure_detector_config = FailureDetectorConfig {
dead_node_grace_period: Duration::from_mins(15),
Expand Down
11 changes: 11 additions & 0 deletions quickwit/quickwit-cluster/src/member.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ pub(crate) const AVAILABILITY_ZONE_KEY: &str = "availability_zone";

pub(crate) const STANDALONE_COMPACTORS_KEY: &str = "standalone_compactors";

pub(crate) const SHARD_SCALING_V2_KEY: &str = "shard_scaling_v2";

pub const INDEXING_CPU_CAPACITY_KEY: &str = "indexing_cpu_capacity";

pub(crate) trait NodeStateExt {
Expand All @@ -56,6 +58,8 @@ pub(crate) trait NodeStateExt {
fn availability_zone(&self) -> Option<AvailabilityZone>;

fn enable_standalone_compactors(&self) -> bool;

fn enable_shard_scaling_v2(&self) -> bool;
}

impl NodeStateExt for NodeState {
Expand Down Expand Up @@ -100,6 +104,10 @@ impl NodeStateExt for NodeState {
fn enable_standalone_compactors(&self) -> bool {
matches!(self.get(STANDALONE_COMPACTORS_KEY), Some(value) if value == "true")
}

fn enable_shard_scaling_v2(&self) -> bool {
matches!(self.get(SHARD_SCALING_V2_KEY), Some(value) if value == "true")
}
}

/// Cluster member.
Expand Down Expand Up @@ -135,6 +143,7 @@ pub struct ClusterMember {
pub availability_zone: Option<AvailabilityZone>,
/// Whether the node was started with standalone compactors enabled.
pub enable_standalone_compactors: bool,
pub enable_shard_scaling_v2: bool,
}

impl ClusterMember {
Expand Down Expand Up @@ -205,6 +214,7 @@ pub(crate) fn build_cluster_member(
let ingester_status = node_state.ingester_status();
let availability_zone = node_state.availability_zone();
let enable_standalone_compactors = node_state.enable_standalone_compactors();
let enable_shard_scaling_v2 = node_state.enable_shard_scaling_v2();

let member = ClusterMember {
node_id: NodeId::from_arc_str(chitchat_id.node_id.clone()),
Expand All @@ -218,6 +228,7 @@ pub(crate) fn build_cluster_member(
ingester_status,
availability_zone,
enable_standalone_compactors,
enable_shard_scaling_v2,
};
Ok(member)
}
Expand Down
4 changes: 4 additions & 0 deletions quickwit/quickwit-cluster/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,10 @@ impl ClusterNode {
pub fn enable_standalone_compactors(&self) -> bool {
self.inner.member.enable_standalone_compactors
}

pub fn enable_shard_scaling_v2(&self) -> bool {
self.inner.member.enable_shard_scaling_v2
}
}

impl std::ops::Deref for ClusterNode {
Expand Down
2 changes: 2 additions & 0 deletions quickwit/quickwit-config/src/node_config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -932,6 +932,8 @@ pub struct NodeConfig {
pub compactor_config: CompactorConfig,
#[serde(skip_serializing)]
pub enable_standalone_compactors: bool,
#[serde(skip_serializing)]
pub enable_shard_scaling_v2: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub docs_clustering_config: Option<DocsClusteringConfig>,
}
Expand Down
6 changes: 6 additions & 0 deletions quickwit/quickwit-config/src/node_config/serialize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,8 @@ struct NodeConfigBuilder {
default_index_root_uri: ConfigValue<Uri, QW_DEFAULT_INDEX_ROOT_URI>,
#[serde(default)]
enable_standalone_compactors: ConfigValue<bool, QW_ENABLE_STANDALONE_COMPACTORS>,
#[serde(default)]
enable_shard_scaling_v2: ConfigValue<bool, QW_ENABLE_SHARD_SCALING_V2>,
#[serde(rename = "rest")]
#[serde(default)]
rest_config_builder: RestConfigBuilder,
Expand Down Expand Up @@ -274,6 +276,7 @@ impl NodeConfigBuilder {
});

let enable_standalone_compactors = self.enable_standalone_compactors.resolve(env_vars)?;
let enable_shard_scaling_v2 = self.enable_shard_scaling_v2.resolve(env_vars)?;
let docs_clustering_config =
DocsClusteringConfigBuilder::build_optional(self.docs_clustering_config, env_vars)?;

Expand Down Expand Up @@ -406,6 +409,7 @@ impl NodeConfigBuilder {
jaeger_config: self.jaeger_config,
compactor_config: self.compactor_config,
enable_standalone_compactors,
enable_shard_scaling_v2,
docs_clustering_config,
};

Expand Down Expand Up @@ -545,6 +549,7 @@ impl Default for NodeConfigBuilder {
metastore_read_replica_uri: default_metastore_read_replica_uri(),
default_index_root_uri: ConfigValue::none(),
enable_standalone_compactors: Default::default(),
enable_shard_scaling_v2: Default::default(),
rest_config_builder: RestConfigBuilder::default(),
health_config_builder: HealthConfigBuilder::default(),
grpc_config: GrpcConfig::default(),
Expand Down Expand Up @@ -708,6 +713,7 @@ pub fn node_config_for_tests_from_ports(
jaeger_config: JaegerConfig::default(),
compactor_config: CompactorConfig::default(),
enable_standalone_compactors: false,
enable_shard_scaling_v2: false,
docs_clustering_config: None,
}
}
Expand Down
5 changes: 3 additions & 2 deletions quickwit/quickwit-config/src/qw_env_vars.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ qw_env_vars!(
QW_DATA_DIR,
QW_DEFAULT_INDEX_ROOT_URI,
QW_DISABLE_DOCS_CLUSTERING,
QW_ENABLE_SHARD_SCALING_V2,
QW_ENABLED_SERVICES,
QW_ENABLE_STANDALONE_COMPACTORS,
QW_EXTRA_CLUSTER_IDS,
Expand Down Expand Up @@ -80,9 +81,9 @@ mod tests {
QW_ENV_VARS.get(&QW_METASTORE_READ_REPLICA_URI).unwrap(),
&"QW_METASTORE_READ_REPLICA_URI"
);
assert_eq!(QW_METASTORE_READ_REPLICA_URI, 16);
assert_eq!(QW_METASTORE_READ_REPLICA_URI, 17);

assert_eq!(QW_ENV_VARS.get(&QW_NODE_ID).unwrap(), &"QW_NODE_ID");
assert_eq!(QW_NODE_ID, 18);
assert_eq!(QW_NODE_ID, 19);
}
}
1 change: 1 addition & 0 deletions quickwit/quickwit-control-plane/src/control_plane.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2169,6 +2169,7 @@ mod tests {
status: IngesterStatus::Retiring,
availability_zone: None,
generation_id: quickwit_cluster::GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);
ingester_pool.insert(
Expand Down
13 changes: 13 additions & 0 deletions quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -395,6 +395,13 @@ impl IngestController {
}
}

pub(crate) fn all_indexers_migrated(&self) -> bool {
self.ingester_pool
.keys_values()
.iter()
.all(|(_node_id, ingester)| ingester.enable_shard_scaling_v2)
}

/// Sends a retain shard request to the given list of ingesters.
///
/// If the request fails, we just log an error.
Expand Down Expand Up @@ -1273,6 +1280,7 @@ mod tests {
status,
availability_zone: availability_zone.map(AvailabilityZone::from),
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
}
}

Expand Down Expand Up @@ -1681,6 +1689,7 @@ mod tests {
status: IngesterStatus::Retiring,
availability_zone: None,
generation_id: quickwit_cluster::GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);
let open_shard_opt =
Expand Down Expand Up @@ -2253,6 +2262,7 @@ mod tests {
status: IngesterStatus::Ready,
availability_zone: Some(AvailabilityZone::from(zone)),
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);
}
Expand Down Expand Up @@ -3013,6 +3023,7 @@ mod tests {
status: IngesterStatus::Ready,
availability_zone: None,
generation_id: quickwit_cluster::GenerationId::from(1u64),
enable_shard_scaling_v2: false,
};
ingester_pool.insert(NodeId::from_str(ingester_id), ingester);
}
Expand All @@ -3031,6 +3042,7 @@ mod tests {
status: IngesterStatus::Retiring,
availability_zone: None,
generation_id: quickwit_cluster::GenerationId::from(1u64),
enable_shard_scaling_v2: false,
};
ingester_pool.insert(NodeId::from_str(ingester_id), ingester);
}
Expand Down Expand Up @@ -3169,6 +3181,7 @@ mod tests {
status: IngesterStatus::Decommissioned,
availability_zone: None,
generation_id: quickwit_cluster::GenerationId::from(1u64),
enable_shard_scaling_v2: false,
};
ingester_pool.insert(NodeId::from_str(ingester_id), ingester);
}
Expand Down
3 changes: 3 additions & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@ pub struct IngesterPoolEntry {
pub status: IngesterStatus,
pub availability_zone: Option<AvailabilityZone>,
pub generation_id: GenerationId,
pub enable_shard_scaling_v2: bool,
}

impl IngesterPoolEntry {
Expand All @@ -80,6 +81,7 @@ impl IngesterPoolEntry {
status: IngesterStatus::Ready,
availability_zone: None,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
}
}

Expand All @@ -90,6 +92,7 @@ impl IngesterPoolEntry {
status: IngesterStatus::Ready,
availability_zone: None,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
}
}
}
Expand Down
7 changes: 7 additions & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1368,6 +1368,7 @@ mod tests {
status: IngesterStatus::Ready,
availability_zone: None,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);

Expand Down Expand Up @@ -1405,6 +1406,7 @@ mod tests {
availability_zone: None,
status: IngesterStatus::Ready,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);

Expand Down Expand Up @@ -1531,6 +1533,7 @@ mod tests {
status: IngesterStatus::Ready,
availability_zone: None,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);

Expand Down Expand Up @@ -1692,6 +1695,7 @@ mod tests {
availability_zone: None,
status: IngesterStatus::Ready,
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
},
);

Expand Down Expand Up @@ -1757,6 +1761,7 @@ mod tests {
NodeId::from_str("test-ingester-0"),
IngesterPoolEntry {
generation_id: GenerationId::from(2u64),
enable_shard_scaling_v2: false,
..IngesterPoolEntry::mocked_ingester()
},
);
Expand Down Expand Up @@ -1912,6 +1917,7 @@ mod tests {
NodeId::from_str("test-ingester-0"),
IngesterPoolEntry {
generation_id: GenerationId::from(3u64),
enable_shard_scaling_v2: false,
..IngesterPoolEntry::mocked_ingester()
},
);
Expand All @@ -1929,6 +1935,7 @@ mod tests {
NodeId::from_str("test-ingester-0"),
IngesterPoolEntry {
generation_id: GenerationId::from(7u64),
enable_shard_scaling_v2: false,
..IngesterPoolEntry::mocked_ingester()
},
);
Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-ingest/src/ingest_v2/routing_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,7 @@ mod tests {
status: IngesterStatus::Ready,
availability_zone: availability_zone.map(AvailabilityZone::from),
generation_id: GenerationId::from(1u64),
enable_shard_scaling_v2: false,
}
}

Expand Down
1 change: 1 addition & 0 deletions quickwit/quickwit-serve/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1280,6 +1280,7 @@ fn build_ingester_insert_change(
status: node.ingester_status,
availability_zone: node.availability_zone(),
generation_id: node.generation_id,
enable_shard_scaling_v2: node.enable_shard_scaling_v2(),
};
Change::Insert(node_id, pool_entry)
}
Expand Down
Loading