Skip to content

Commit f2b97a5

Browse files
committed
error out on node id hash collision
1 parent 6d7383e commit f2b97a5

3 files changed

Lines changed: 63 additions & 39 deletions

File tree

rust/operator-binary/src/config/node_id_hasher.rs

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,8 @@ pub fn node_id_hash32_offset(rolegroup_ref: &RoleGroupRef<KafkaCluster>) -> u32
99
rolegroup = rolegroup_ref.role_group
1010
));
1111
let range = hash & 0x0000FFFF;
12-
// unsigned in kafka
13-
let offset = range * 0x00007FFF;
14-
offset
12+
// Kafka uses signed integer
13+
range * 0x00007FFF
1514
}
1615

1716
/// Simple FNV-1a hash impl

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

Lines changed: 59 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ pub mod role;
66
pub mod security;
77
pub mod tls;
88

9-
use std::collections::BTreeMap;
9+
use std::collections::{BTreeMap, HashMap};
1010

1111
use authentication::KafkaAuthentication;
1212
use serde::{Deserialize, Serialize};
@@ -60,10 +60,10 @@ pub const STACKABLE_KERBEROS_KRB5_PATH: &str = "/stackable/kerberos/krb5.conf";
6060

6161
#[derive(Snafu, Debug)]
6262
pub enum Error {
63-
#[snafu(display("the Kafka role [{role}] is missing from spec"))]
63+
#[snafu(display("The Kafka role [{role}] is missing from spec"))]
6464
MissingRole { role: String },
6565

66-
#[snafu(display("object has no namespace associated"))]
66+
#[snafu(display("Object has no namespace associated"))]
6767
NoNamespace,
6868

6969
#[snafu(display(
@@ -75,6 +75,15 @@ pub enum Error {
7575
"Kraft controller (`spec.controller`) and ZooKeeper (`spec.clusterConfig.zookeeperConfigMapName`) are configured. Please only choose one"
7676
))]
7777
KraftAndZookeeperConfigured,
78+
79+
#[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."
81+
))]
82+
KafkaNodeIdHashCollision {
83+
role: KafkaRole,
84+
rolegroup: String,
85+
colliding_rolegroup: String,
86+
},
7887
}
7988

8089
#[versioned(
@@ -241,49 +250,65 @@ impl v1alpha1::KafkaCluster {
241250
kafka_role: &KafkaRole,
242251
cluster_info: &KubernetesClusterInfo,
243252
) -> Result<Vec<KafkaPodDescriptor>, Error> {
244-
let ns = self.metadata.namespace.clone().context(NoNamespaceSnafu)?;
253+
let namespace = self.metadata.namespace.clone().context(NoNamespaceSnafu)?;
254+
let rolegroup_replicas = self.extract_rolegroup_replicas(kafka_role)?;
255+
let mut pod_descriptors = Vec::new();
256+
let mut seen_hashes = HashMap::<u32, String>::new();
257+
258+
for (rolegroup, replicas) in rolegroup_replicas {
259+
let rolegroup_ref = self.rolegroup_ref(kafka_role, &rolegroup);
260+
let node_id_hash_offset = node_id_hash32_offset(&rolegroup_ref);
261+
262+
match seen_hashes.get(&node_id_hash_offset) {
263+
Some(colliding_rolegroup) => {
264+
return KafkaNodeIdHashCollisionSnafu {
265+
role: kafka_role.clone(),
266+
rolegroup: rolegroup.clone(),
267+
colliding_rolegroup: colliding_rolegroup.clone(),
268+
}
269+
.fail();
270+
}
271+
None => seen_hashes.insert(node_id_hash_offset, rolegroup),
272+
};
273+
274+
for replica in 0..replicas {
275+
pod_descriptors.push(KafkaPodDescriptor {
276+
namespace: namespace.clone(),
277+
role_group_service_name: rolegroup_ref.object_name(),
278+
replica,
279+
cluster_domain: cluster_info.cluster_domain.clone(),
280+
node_id: node_id_hash_offset + u32::from(replica),
281+
});
282+
}
283+
}
284+
285+
Ok(pod_descriptors)
286+
}
287+
288+
fn extract_rolegroup_replicas(
289+
&self,
290+
kafka_role: &KafkaRole,
291+
) -> Result<BTreeMap<String, u16>, Error> {
245292
Ok(match kafka_role {
246293
KafkaRole::Broker => self
247294
.broker_role()
248295
.iter()
249296
.flat_map(|role| &role.role_groups)
250-
// Order rolegroups consistently, to avoid spurious downstream rewrites
251-
.collect::<BTreeMap<_, _>>()
252-
.into_iter()
253-
.flat_map(move |(rolegroup_name, rolegroup)| {
254-
let rolegroup_ref = self.rolegroup_ref(kafka_role, rolegroup_name);
255-
let ns = ns.clone();
256-
(0..rolegroup.replicas.unwrap_or(0)).map(move |i| KafkaPodDescriptor {
257-
namespace: ns.clone(),
258-
role_group_service_name: rolegroup_ref.object_name(),
259-
replica: i,
260-
cluster_domain: cluster_info.cluster_domain.clone(),
261-
// TODO: check for hash collisions?
262-
node_id: node_id_hash32_offset(&rolegroup_ref) + u32::from(i),
263-
})
297+
.flat_map(|(rolegroup_name, rolegroup)| {
298+
std::iter::once((rolegroup_name.to_string(), rolegroup.replicas.unwrap_or(0)))
264299
})
265-
.collect(),
300+
// Order rolegroups consistently, to avoid spurious downstream rewrites
301+
.collect::<BTreeMap<_, _>>(),
266302

267303
KafkaRole::Controller => self
268304
.controller_role()
269305
.iter()
270306
.flat_map(|role| &role.role_groups)
271-
// Order rolegroups consistently, to avoid spurious downstream rewrites
272-
.collect::<BTreeMap<_, _>>()
273-
.into_iter()
274-
.flat_map(move |(rolegroup_name, rolegroup)| {
275-
let rolegroup_ref = self.rolegroup_ref(kafka_role, rolegroup_name);
276-
let ns = ns.clone();
277-
(0..rolegroup.replicas.unwrap_or(0)).map(move |i: u16| KafkaPodDescriptor {
278-
namespace: ns.clone(),
279-
role_group_service_name: rolegroup_ref.object_name(),
280-
replica: i,
281-
cluster_domain: cluster_info.cluster_domain.clone(),
282-
// TODO: check for hash collisions?
283-
node_id: node_id_hash32_offset(&rolegroup_ref) + u32::from(i),
284-
})
307+
.flat_map(|(rolegroup_name, rolegroup)| {
308+
std::iter::once((rolegroup_name.to_string(), rolegroup.replicas.unwrap_or(0)))
285309
})
286-
.collect(),
310+
// Order rolegroups consistently, to avoid spurious downstream rewrites
311+
.collect::<BTreeMap<_, _>>(),
287312
})
288313
}
289314
}

rust/operator-binary/src/resource/statefulset.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -333,7 +333,7 @@ pub fn build_broker_rolegroup_statefulset(
333333
)
334334
.add_env_var(
335335
KAFKA_NODE_ID_OFFSET,
336-
node_id_hash32_offset(&rolegroup_ref).to_string(),
336+
node_id_hash32_offset(rolegroup_ref).to_string(),
337337
)
338338
.add_env_vars(env)
339339
.add_container_ports(container_ports(kafka_security))
@@ -675,7 +675,7 @@ pub fn build_controller_rolegroup_statefulset(
675675
)
676676
.add_env_var(
677677
KAFKA_NODE_ID_OFFSET,
678-
node_id_hash32_offset(&rolegroup_ref).to_string(),
678+
node_id_hash32_offset(rolegroup_ref).to_string(),
679679
)
680680
.add_env_vars(env)
681681
.add_container_ports(container_ports(kafka_security))

0 commit comments

Comments
 (0)