@@ -10,13 +10,8 @@ use crate::{
1010 crd:: {
1111 KafkaPodDescriptor , STACKABLE_CONFIG_DIR , STACKABLE_KERBEROS_KRB5_PATH ,
1212 STACKABLE_LISTENER_BOOTSTRAP_DIR , STACKABLE_LISTENER_BROKER_DIR ,
13- listener:: { KafkaListenerConfig , KafkaListenerName , node_address_cmd} ,
14- role:: {
15- KAFKA_ADVERTISED_LISTENERS , KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS ,
16- KAFKA_CONTROLLER_QUORUM_VOTERS , KAFKA_LISTENER_SECURITY_PROTOCOL_MAP , KAFKA_LISTENERS ,
17- KAFKA_NODE_ID , KAFKA_NODE_ID_OFFSET , KafkaRole , broker:: BROKER_PROPERTIES_FILE ,
18- controller:: CONTROLLER_PROPERTIES_FILE ,
19- } ,
13+ listener:: { KafkaListenerName , node_address_cmd} ,
14+ role:: { KafkaRole , broker:: BROKER_PROPERTIES_FILE , controller:: CONTROLLER_PROPERTIES_FILE } ,
2015 security:: KafkaTlsSecurity ,
2116 v1alpha1,
2217 } ,
@@ -28,8 +23,6 @@ pub fn broker_kafka_container_commands(
2823 kafka : & v1alpha1:: KafkaCluster ,
2924 cluster_id : & str ,
3025 controller_descriptors : Vec < KafkaPodDescriptor > ,
31- kafka_listeners : & KafkaListenerConfig ,
32- opa_connect_string : Option < & str > ,
3326 kafka_security : & KafkaTlsSecurity ,
3427 product_version : & str ,
3528) -> String {
@@ -51,26 +44,17 @@ pub fn broker_kafka_container_commands(
5144 true => format!( "export KERBEROS_REALM=$(grep -oP 'default_realm = \\ K.*' {STACKABLE_KERBEROS_KRB5_PATH})" ) ,
5245 false => "" . to_string( ) ,
5346 } ,
54- broker_start_command = broker_start_command( kafka, cluster_id, controller_descriptors, kafka_listeners , opa_connect_string , kafka_security, product_version) ,
47+ broker_start_command = broker_start_command( kafka, cluster_id, controller_descriptors, kafka_security, product_version) ,
5548 }
5649}
5750
5851fn broker_start_command (
5952 kafka : & v1alpha1:: KafkaCluster ,
6053 cluster_id : & str ,
6154 controller_descriptors : Vec < KafkaPodDescriptor > ,
62- kafka_listeners : & KafkaListenerConfig ,
63- opa_connect_string : Option < & str > ,
6455 kafka_security : & KafkaTlsSecurity ,
6556 product_version : & str ,
6657) -> String {
67- let opa_config = match opa_connect_string {
68- None => "" . to_string ( ) ,
69- Some ( opa_connect_string) => {
70- format ! ( " --override \" opa.authorizer.url={opa_connect_string}\" " )
71- }
72- } ;
73-
7458 let jaas_config = match kafka_security. has_kerberos_enabled ( ) {
7559 true => {
7660 formatdoc ! { "
@@ -89,50 +73,35 @@ fn broker_start_command(
8973
9074 let client_port = kafka_security. client_port ( ) ;
9175
92- // TODO: The properties file from the configmap is copied to the /tmp folder and appended with dynamic properties
9376 // This should be improved:
9477 // - mount emptyDir as readWriteConfig
95- // - use config-utils for proper replacements?
96- // - should we print the adapted properties file at startup?
9778 if kafka. is_controller_configured ( ) {
9879 formatdoc ! { "
99- export REPLICA_ID=$(echo \" $POD_NAME\" | grep -oE '[0-9]+$')
80+ POD_INDEX=$(echo \" $POD_NAME\" | grep -oE '[0-9]+$')
81+ export REPLICA_ID=$((POD_INDEX+NODE_ID_OFFSET))
82+
10083 cp {config_dir}/{properties_file} /tmp/{properties_file}
10184
102- echo \" {KAFKA_NODE_ID}=$((REPLICA_ID + ${KAFKA_NODE_ID_OFFSET}))\" >> /tmp/{properties_file}
103- echo \" {KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS}={bootstrap_servers}\" >> /tmp/{properties_file}
104- echo \" {KAFKA_LISTENERS}={listeners}\" >> /tmp/{properties_file}
105- echo \" {KAFKA_ADVERTISED_LISTENERS}={advertised_listeners}\" >> /tmp/{properties_file}
106- echo \" {KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" >> /tmp/{properties_file}
107- echo \" {KAFKA_CONTROLLER_QUORUM_VOTERS}={controller_quorum_voters}\" >> /tmp/{properties_file}
85+ config-utils template /tmp/{properties_file}
10886
10987 bin/kafka-storage.sh format --cluster-id {cluster_id} --config /tmp/{properties_file} --ignore-formatted {initial_controller_command}
110- bin/kafka-server-start.sh /tmp/{properties_file} {opa_config}{ jaas_config} &
88+ bin/kafka-server-start.sh /tmp/{properties_file} {jaas_config} &
11189 " ,
11290 config_dir = STACKABLE_CONFIG_DIR ,
11391 properties_file = BROKER_PROPERTIES_FILE ,
114- bootstrap_servers = to_bootstrap_servers( & controller_descriptors, client_port) ,
115- listeners = kafka_listeners. listeners( ) ,
116- advertised_listeners = kafka_listeners. advertised_listeners( ) ,
117- listener_security_protocol_map = kafka_listeners. listener_security_protocol_map( ) ,
118- controller_quorum_voters = to_quorum_voters( & controller_descriptors, client_port) ,
11992 initial_controller_command = initial_controllers_command( & controller_descriptors, product_version, client_port) ,
12093 }
12194 } else {
12295 formatdoc ! { "
123- bin/kafka-server-start.sh {config_dir}/{properties_file} \
124- --override \" zookeeper.connect=$ZOOKEEPER\" \
125- --override \" {KAFKA_LISTENERS}={listeners}\" \
126- --override \" {KAFKA_ADVERTISED_LISTENERS}={advertised_listeners}\" \
127- --override \" {KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" \
128- {opa_config} \
96+ cp {config_dir}/{properties_file} /tmp/{properties_file}
97+
98+ config-utils template /tmp/{properties_file}
99+
100+ bin/kafka-server-start.sh /tmp/{properties_file} \
129101 {jaas_config} \
130102 &",
131103 config_dir = STACKABLE_CONFIG_DIR ,
132104 properties_file = BROKER_PROPERTIES_FILE ,
133- listeners = kafka_listeners. listeners( ) ,
134- advertised_listeners = kafka_listeners. advertised_listeners( ) ,
135- listener_security_protocol_map = kafka_listeners. listener_security_protocol_map( ) ,
136105 }
137106 }
138107}
@@ -182,7 +151,6 @@ wait_for_termination()
182151pub fn controller_kafka_container_command (
183152 cluster_id : & str ,
184153 controller_descriptors : Vec < KafkaPodDescriptor > ,
185- kafka_listeners : & KafkaListenerConfig ,
186154 kafka_security : & KafkaTlsSecurity ,
187155 product_version : & str ,
188156) -> String {
@@ -199,14 +167,12 @@ pub fn controller_kafka_container_command(
199167 prepare_signal_handlers
200168 containerdebug --output={STACKABLE_LOG_DIR}/containerdebug-state.json --loop &
201169
202- export REPLICA_ID=$(echo \" $POD_NAME\" | grep -oE '[0-9]+$')
170+ POD_INDEX=$(echo \" $POD_NAME\" | grep -oE '[0-9]+$')
171+ export REPLICA_ID=$((POD_INDEX+NODE_ID_OFFSET))
172+
203173 cp {config_dir}/{properties_file} /tmp/{properties_file}
204174
205- echo \" {KAFKA_NODE_ID}=$((REPLICA_ID + ${KAFKA_NODE_ID_OFFSET}))\" >> /tmp/{properties_file}
206- echo \" {KAFKA_CONTROLLER_QUORUM_BOOTSTRAP_SERVERS}={bootstrap_servers}\" >> /tmp/{properties_file}
207- echo \" {KAFKA_LISTENERS}={listeners}\" >> /tmp/{properties_file}
208- echo \" {KAFKA_LISTENER_SECURITY_PROTOCOL_MAP}={listener_security_protocol_map}\" >> /tmp/{properties_file}
209- echo \" {KAFKA_CONTROLLER_QUORUM_VOTERS}={controller_quorum_voters}\" >> /tmp/{properties_file}
175+ config-utils template /tmp/{properties_file}
210176
211177 bin/kafka-storage.sh format --cluster-id {cluster_id} --config /tmp/{properties_file} --ignore-formatted {initial_controller_command}
212178 bin/kafka-server-start.sh /tmp/{properties_file} &
@@ -217,29 +183,11 @@ pub fn controller_kafka_container_command(
217183 remove_vector_shutdown_file_command = remove_vector_shutdown_file_command( STACKABLE_LOG_DIR ) ,
218184 config_dir = STACKABLE_CONFIG_DIR ,
219185 properties_file = CONTROLLER_PROPERTIES_FILE ,
220- bootstrap_servers = to_bootstrap_servers( & controller_descriptors, client_port) ,
221- listeners = to_listeners( client_port) ,
222- listener_security_protocol_map = to_listener_security_protocol_map( kafka_listeners) ,
223186 initial_controller_command = initial_controllers_command( & controller_descriptors, product_version, client_port) ,
224- controller_quorum_voters = to_quorum_voters( & controller_descriptors, client_port) ,
225187 create_vector_shutdown_file_command = create_vector_shutdown_file_command( STACKABLE_LOG_DIR )
226188 }
227189}
228190
229- fn to_listeners ( port : u16 ) -> String {
230- // The environment variables are set in the statefulset of the controller
231- format ! (
232- "{listener_name}://$POD_NAME.$ROLEGROUP_HEADLESS_SERVICE_NAME.$NAMESPACE.svc.$CLUSTER_DOMAIN:{port}" ,
233- listener_name = KafkaListenerName :: Controller
234- )
235- }
236-
237- fn to_listener_security_protocol_map ( kafka_listeners : & KafkaListenerConfig ) -> String {
238- kafka_listeners
239- . listener_security_protocol_map_for_listener ( & KafkaListenerName :: Controller )
240- . unwrap_or ( "" . to_string ( ) )
241- }
242-
243191fn to_initial_controllers ( controller_descriptors : & [ KafkaPodDescriptor ] , port : u16 ) -> String {
244192 controller_descriptors
245193 . iter ( )
@@ -248,23 +196,6 @@ fn to_initial_controllers(controller_descriptors: &[KafkaPodDescriptor], port: u
248196 . join ( "," )
249197}
250198
251- // TODO: This can be removed once 3.7.2 is removed. Used in command.rs.
252- fn to_quorum_voters ( controller_descriptors : & [ KafkaPodDescriptor ] , port : u16 ) -> String {
253- controller_descriptors
254- . iter ( )
255- . map ( |desc| desc. as_quorum_voter ( port) )
256- . collect :: < Vec < String > > ( )
257- . join ( "," )
258- }
259-
260- fn to_bootstrap_servers ( controller_descriptors : & [ KafkaPodDescriptor ] , port : u16 ) -> String {
261- controller_descriptors
262- . iter ( )
263- . map ( |desc| format ! ( "{fqdn}:{port}" , fqdn = desc. fqdn( ) ) )
264- . collect :: < Vec < String > > ( )
265- . join ( "," )
266- }
267-
268199fn initial_controllers_command (
269200 controller_descriptors : & [ KafkaPodDescriptor ] ,
270201 product_version : & str ,
0 commit comments