Skip to content

Commit 7cc2e2f

Browse files
committed
remove second metrics service, consolidate ports
1 parent d572fa3 commit 7cc2e2f

5 files changed

Lines changed: 51 additions & 117 deletions

File tree

rust/operator-binary/src/container.rs

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -488,9 +488,7 @@ impl ContainerConfig {
488488
)?)
489489
.add_volume_mounts(self.volume_mounts(hdfs, merged_config, labels)?)
490490
.context(AddVolumeMountSnafu)?
491-
.add_container_ports(self.container_ports(hdfs))
492-
// TODO: This currently adds the metrics port also to the zkfc containers, not needed there?
493-
.add_container_port(SERVICE_PORT_NAME_METRICS, hdfs.metrics_port(role).into());
491+
.add_container_ports(self.container_ports(hdfs));
494492

495493
if let Some(resources) = resources {
496494
cb.resources(resources);
@@ -1251,16 +1249,18 @@ wait_for_termination $!
12511249
/// Container ports for the main containers namenode, datanode and journalnode.
12521250
fn container_ports(&self, hdfs: &v1alpha1::HdfsCluster) -> Vec<ContainerPort> {
12531251
match self {
1254-
ContainerConfig::Hdfs { role, .. } => hdfs
1255-
.ports(role)
1256-
.into_iter()
1257-
.map(|(name, value)| ContainerPort {
1258-
name: Some(name),
1259-
container_port: i32::from(value),
1260-
protocol: Some("TCP".to_string()),
1261-
..ContainerPort::default()
1262-
})
1263-
.collect(),
1252+
ContainerConfig::Hdfs { role, .. } => {
1253+
// data ports
1254+
hdfs.hdfs_main_container_ports(role)
1255+
.into_iter()
1256+
.map(|(name, value)| ContainerPort {
1257+
name: Some(name),
1258+
container_port: i32::from(value),
1259+
protocol: Some("TCP".to_string()),
1260+
..ContainerPort::default()
1261+
})
1262+
.collect()
1263+
}
12641264
_ => {
12651265
vec![]
12661266
}

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ pub const SERVICE_PORT_NAME_HTTP: &str = "http";
2020
pub const SERVICE_PORT_NAME_HTTPS: &str = "https";
2121
pub const SERVICE_PORT_NAME_DATA: &str = "data";
2222
pub const SERVICE_PORT_NAME_METRICS: &str = "metrics";
23+
pub const SERVICE_PORT_NAME_JMX_METRICS: &str = "jmx-metrics";
2324

2425
pub const DEFAULT_LISTENER_CLASS: &str = "cluster-internal";
2526

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

Lines changed: 34 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,8 @@ use crate::crd::{
6565
DEFAULT_NAME_NODE_RPC_PORT, DFS_REPLICATION, HADOOP_POLICY_XML, HDFS_SITE_XML,
6666
JVM_SECURITY_PROPERTIES_FILE, LISTENER_VOLUME_NAME, SERVICE_PORT_NAME_DATA,
6767
SERVICE_PORT_NAME_HTTP, SERVICE_PORT_NAME_HTTPS, SERVICE_PORT_NAME_IPC,
68-
SERVICE_PORT_NAME_METRICS, SERVICE_PORT_NAME_RPC, SSL_CLIENT_XML, SSL_SERVER_XML,
68+
SERVICE_PORT_NAME_JMX_METRICS, SERVICE_PORT_NAME_METRICS, SERVICE_PORT_NAME_RPC,
69+
SSL_CLIENT_XML, SSL_SERVER_XML,
6970
},
7071
security::{AuthenticationConfig, KerberosConfig},
7172
storage::{
@@ -390,10 +391,10 @@ impl v1alpha1::HdfsCluster {
390391
let ns = ns.clone();
391392
(0..*replicas).map(move |i| HdfsPodRef {
392393
namespace: ns.clone(),
393-
role_group_service_name: rolegroup_ref.object_name(),
394+
role_group_service_name: rolegroup_ref.rolegroup_headless_service_name(),
394395
pod_name: format!("{}-{}", rolegroup_ref.object_name(), i),
395396
ports: self
396-
.ports(role)
397+
.data_ports(role)
397398
.iter()
398399
.map(|(n, p)| (n.clone(), *p))
399400
.collect(),
@@ -671,7 +672,7 @@ impl v1alpha1::HdfsCluster {
671672
}
672673

673674
/// Returns required port name and port number tuples depending on the role.
674-
pub fn ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
675+
pub fn data_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
675676
match role {
676677
HdfsNodeRole::Name => vec![
677678
(
@@ -731,26 +732,26 @@ impl v1alpha1::HdfsCluster {
731732
}
732733
}
733734

734-
/// Returns required metrics port name and metrics port number tuples depending on the role.
735-
pub fn metrics_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
735+
/// Deprecated required JMX metrics port name and metrics port number tuples depending on the role.
736+
pub fn jmx_metrics_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
736737
match role {
737738
HdfsNodeRole::Name => vec![(
738-
String::from(SERVICE_PORT_NAME_METRICS),
739+
String::from(SERVICE_PORT_NAME_JMX_METRICS),
739740
DEFAULT_NAME_NODE_METRICS_PORT,
740741
)],
741742
HdfsNodeRole::Data => vec![(
742-
String::from(SERVICE_PORT_NAME_METRICS),
743+
String::from(SERVICE_PORT_NAME_JMX_METRICS),
743744
DEFAULT_DATA_NODE_METRICS_PORT,
744745
)],
745746
HdfsNodeRole::Journal => vec![(
746-
String::from(SERVICE_PORT_NAME_METRICS),
747+
String::from(SERVICE_PORT_NAME_JMX_METRICS),
747748
DEFAULT_JOURNAL_NODE_METRICS_PORT,
748749
)],
749750
}
750751
}
751752

752-
/// Returns required metrics port name and native metrics port number tuples depending on the role.
753-
pub fn native_metrics_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
753+
/// Returns required metrics port name and metrics port number tuples depending on the role and security settings.
754+
pub fn metrics_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
754755
match role {
755756
HdfsNodeRole::Name => vec![if self.has_https_enabled() {
756757
(
@@ -821,6 +822,28 @@ impl v1alpha1::HdfsCluster {
821822
}
822823
}
823824
}
825+
826+
pub fn metrics_service_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
827+
let mut metrics_service_ports = vec![];
828+
// "native" ports
829+
metrics_service_ports.extend(self.metrics_ports(role));
830+
metrics_service_ports.extend(self.jmx_metrics_ports(role));
831+
metrics_service_ports
832+
}
833+
834+
pub fn headless_service_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
835+
let mut headless_service_ports = vec![];
836+
headless_service_ports.extend(self.data_ports(role));
837+
headless_service_ports
838+
}
839+
840+
pub fn hdfs_main_container_ports(&self, role: &HdfsNodeRole) -> Vec<(String, u16)> {
841+
let mut main_container_ports = vec![];
842+
main_container_ports.extend(self.data_ports(role));
843+
// TODO: This will be exposed in the listener if added to container ports?
844+
// main_container_ports.extend(self.jmx_metrics_ports(role));
845+
main_container_ports
846+
}
824847
}
825848

826849
#[derive(Clone, Debug, Deserialize, Eq, Hash, JsonSchema, PartialEq, Serialize)]

rust/operator-binary/src/hdfs_controller.rs

Lines changed: 1 addition & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -69,10 +69,7 @@ use crate::{
6969
},
7070
product_logging::extend_role_group_config_map,
7171
security::{self, kerberos, opa::HdfsOpaConfig},
72-
service::{
73-
self, rolegroup_headless_service, rolegroup_metrics_service,
74-
rolegroup_native_metrics_service,
75-
},
72+
service::{self, rolegroup_headless_service, rolegroup_metrics_service},
7673
};
7774

7875
pub const RESOURCE_MANAGER_HDFS_CONTROLLER: &str = "hdfs-operator-hdfs-controller";
@@ -399,13 +396,6 @@ pub async fn reconcile_hdfs(
399396
let rg_metrics_service =
400397
rolegroup_metrics_service(hdfs, &role, &rolegroup_ref, &resolved_product_image)
401398
.context(BuildServiceSnafu)?;
402-
let rg_native_metrics_service = rolegroup_native_metrics_service(
403-
hdfs,
404-
&role,
405-
&rolegroup_ref,
406-
&resolved_product_image,
407-
)
408-
.context(BuildServiceSnafu)?;
409399

410400
// We need to split the creation and the usage of the "metadata" variable in two statements.
411401
// to avoid the compiler error "E0716 (temporary value dropped while borrowed)".
@@ -453,7 +443,6 @@ pub async fn reconcile_hdfs(
453443

454444
let rg_service_name = rg_service.name_any();
455445
let rg_metrics_service_name = rg_metrics_service.name_any();
456-
let rg_native_metrics_service_name = rg_native_metrics_service.name_any();
457446

458447
cluster_resources
459448
.add(client, rg_service)
@@ -467,12 +456,6 @@ pub async fn reconcile_hdfs(
467456
.with_context(|_| ApplyRoleGroupServiceSnafu {
468457
name: rg_metrics_service_name,
469458
})?;
470-
cluster_resources
471-
.add(client, rg_native_metrics_service)
472-
.await
473-
.with_context(|_| ApplyRoleGroupServiceSnafu {
474-
name: rg_native_metrics_service_name,
475-
})?;
476459
let rg_configmap_name = rg_configmap.name_any();
477460
cluster_resources
478461
.add(client, rg_configmap.clone())

rust/operator-binary/src/service.rs

Lines changed: 2 additions & 75 deletions
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ pub(crate) fn rolegroup_headless_service(
6767
type_: Some("ClusterIP".to_string()),
6868
cluster_ip: Some("None".to_string()),
6969
ports: Some(
70-
hdfs.ports(role)
70+
hdfs.headless_service_ports(role)
7171
.into_iter()
7272
.map(|(name, value)| ServicePort {
7373
name: Some(name),
@@ -106,7 +106,7 @@ pub(crate) fn rolegroup_metrics_service(
106106
type_: Some("ClusterIP".to_string()),
107107
cluster_ip: Some("None".to_string()),
108108
ports: Some(
109-
hdfs.metrics_ports(role)
109+
hdfs.metrics_service_ports(role)
110110
.into_iter()
111111
.map(|(name, value)| ServicePort {
112112
name: Some(name),
@@ -145,79 +145,6 @@ pub(crate) fn rolegroup_metrics_service(
145145
Label::try_from(("prometheus.io/scrape", "true"))
146146
.context(BuildPrometheusLabelSnafu)?,
147147
)
148-
.with_annotations(
149-
Annotations::try_from([
150-
("prometheus.io/path".to_owned(), "/metrics".to_owned()),
151-
(
152-
"prometheus.io/port".to_owned(),
153-
hdfs.metrics_port(role).to_string(),
154-
),
155-
("prometheus.io/scheme".to_owned(), "http".to_owned()),
156-
("prometheus.io/scrape".to_owned(), "true".to_owned()),
157-
])
158-
.expect("should be valid annotations"),
159-
)
160-
.build(),
161-
spec: Some(service_spec),
162-
status: None,
163-
})
164-
}
165-
166-
pub(crate) fn rolegroup_native_metrics_service(
167-
hdfs: &v1alpha1::HdfsCluster,
168-
role: &HdfsNodeRole,
169-
rolegroup_ref: &RoleGroupRef<v1alpha1::HdfsCluster>,
170-
resolved_product_image: &ResolvedProductImage,
171-
) -> Result<Service, Error> {
172-
tracing::info!("Setting up native metrics Service for {:?}", rolegroup_ref);
173-
174-
let service_spec = ServiceSpec {
175-
// Internal communication does not need to be exposed
176-
type_: Some("ClusterIP".to_string()),
177-
cluster_ip: Some("None".to_string()),
178-
ports: Some(
179-
hdfs.native_metrics_ports(role)
180-
.into_iter()
181-
.map(|(name, value)| ServicePort {
182-
name: Some(name),
183-
port: i32::from(value),
184-
protocol: Some("TCP".to_string()),
185-
..ServicePort::default()
186-
})
187-
.collect(),
188-
),
189-
selector: Some(
190-
hdfs.rolegroup_selector_labels(rolegroup_ref)
191-
.context(RoleGroupSelectorLabelsSnafu)?
192-
.into(),
193-
),
194-
publish_not_ready_addresses: Some(true),
195-
..ServiceSpec::default()
196-
};
197-
198-
Ok(Service {
199-
metadata: ObjectMetaBuilder::new()
200-
.name_and_namespace(hdfs)
201-
.name(format!(
202-
"{name}-native-metrics",
203-
name = rolegroup_ref.object_name()
204-
))
205-
.ownerreference_from_resource(hdfs, None, Some(true))
206-
.with_context(|_| ObjectMissingMetadataForOwnerRefSnafu {
207-
obj_ref: ObjectRef::from_obj(hdfs),
208-
})?
209-
.with_recommended_labels(build_recommended_labels(
210-
hdfs,
211-
RESOURCE_MANAGER_HDFS_CONTROLLER,
212-
&resolved_product_image.app_version_label_value,
213-
&rolegroup_ref.role,
214-
&rolegroup_ref.role_group,
215-
))
216-
.context(ObjectMetaSnafu)?
217-
.with_label(
218-
Label::try_from(("prometheus.io/scrape", "true"))
219-
.context(BuildPrometheusLabelSnafu)?,
220-
)
221148
.with_annotations(
222149
Annotations::try_from([
223150
("prometheus.io/path".to_owned(), "/prom".to_owned()),

0 commit comments

Comments
 (0)