Skip to content

Commit 43a1759

Browse files
committed
implement node.id hashing check for all roles
1 parent 91659bb commit 43a1759

2 files changed

Lines changed: 39 additions & 30 deletions

File tree

rust/operator-binary/src/crd/mod.rs

Lines changed: 37 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -77,11 +77,12 @@ pub enum Error {
7777
KraftAndZookeeperConfigured,
7878

7979
#[snafu(display(
80-
"Could not calculate ({role}) 'node.id' hash offset for rolegroup '{rolegroup}' which collides with rolegroup '{colliding_rolegroup}'. Please try to rename one of the rolegroups."
80+
"Could not calculate 'node.id' hash offset for role '{role}' and rolegroup '{rolegroup}' which collides with role '{coliding_role}' and rolegroup '{colliding_rolegroup}'. Please try to rename one of the rolegroups."
8181
))]
8282
KafkaNodeIdHashCollision {
8383
role: KafkaRole,
8484
rolegroup: String,
85+
coliding_role: KafkaRole,
8586
colliding_rolegroup: String,
8687
},
8788
}
@@ -258,41 +259,49 @@ impl v1alpha1::KafkaCluster {
258259
///
259260
/// We try to predict the pods here rather than looking at the current cluster state in order to
260261
/// avoid instance churn.
261-
// TODO: this currently only checks within each role, node.id must be unique for all brokers and controllers
262262
pub fn pod_descriptors(
263263
&self,
264-
kafka_role: &KafkaRole,
264+
requested_kafka_role: &KafkaRole,
265265
cluster_info: &KubernetesClusterInfo,
266266
) -> Result<Vec<KafkaPodDescriptor>, Error> {
267267
let namespace = self.metadata.namespace.clone().context(NoNamespaceSnafu)?;
268-
let rolegroup_replicas = self.extract_rolegroup_replicas(kafka_role)?;
269268
let mut pod_descriptors = Vec::new();
270-
let mut seen_hashes = HashMap::<u32, String>::new();
271-
272-
for (rolegroup, replicas) in rolegroup_replicas {
273-
let rolegroup_ref = self.rolegroup_ref(kafka_role, &rolegroup);
274-
let node_id_hash_offset = node_id_hash32_offset(&rolegroup_ref);
275-
276-
match seen_hashes.get(&node_id_hash_offset) {
277-
Some(colliding_rolegroup) => {
278-
return KafkaNodeIdHashCollisionSnafu {
279-
role: kafka_role.clone(),
280-
rolegroup: rolegroup.clone(),
281-
colliding_rolegroup: colliding_rolegroup.clone(),
269+
let mut seen_hashes = HashMap::<u32, (KafkaRole, String)>::new();
270+
271+
for current_role in KafkaRole::roles() {
272+
let rolegroup_replicas = self.extract_rolegroup_replicas(&current_role)?;
273+
for (rolegroup, replicas) in rolegroup_replicas {
274+
let rolegroup_ref = self.rolegroup_ref(&current_role, &rolegroup);
275+
let node_id_hash_offset = node_id_hash32_offset(&rolegroup_ref);
276+
277+
// check collisions
278+
match seen_hashes.get(&node_id_hash_offset) {
279+
Some((coliding_role, coliding_rolegroup)) => {
280+
return KafkaNodeIdHashCollisionSnafu {
281+
role: current_role.clone(),
282+
rolegroup: rolegroup.clone(),
283+
coliding_role: coliding_role.clone(),
284+
colliding_rolegroup: coliding_rolegroup.to_string(),
285+
}
286+
.fail();
287+
}
288+
None => {
289+
seen_hashes.insert(node_id_hash_offset, (current_role.clone(), rolegroup))
290+
}
291+
};
292+
293+
// only return descriptors for selected role
294+
if current_role == *requested_kafka_role {
295+
for replica in 0..replicas {
296+
pod_descriptors.push(KafkaPodDescriptor {
297+
namespace: namespace.clone(),
298+
role_group_service_name: rolegroup_ref.object_name(),
299+
replica,
300+
cluster_domain: cluster_info.cluster_domain.clone(),
301+
node_id: node_id_hash_offset + u32::from(replica),
302+
});
282303
}
283-
.fail();
284304
}
285-
None => seen_hashes.insert(node_id_hash_offset, rolegroup),
286-
};
287-
288-
for replica in 0..replicas {
289-
pod_descriptors.push(KafkaPodDescriptor {
290-
namespace: namespace.clone(),
291-
role_group_service_name: rolegroup_ref.object_name(),
292-
replica,
293-
cluster_domain: cluster_info.cluster_domain.clone(),
294-
node_id: node_id_hash_offset + u32::from(replica),
295-
});
296305
}
297306
}
298307

rust/operator-binary/src/crd/role/mod.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -110,10 +110,10 @@ pub enum KafkaRole {
110110

111111
impl KafkaRole {
112112
/// Return all available roles
113-
pub fn roles() -> Vec<String> {
113+
pub fn roles() -> Vec<KafkaRole> {
114114
let mut roles = vec![];
115115
for role in Self::iter() {
116-
roles.push(role.to_string())
116+
roles.push(role)
117117
}
118118
roles
119119
}

0 commit comments

Comments
 (0)