Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -11180,6 +11180,71 @@ spec:
type: string
type: object
type: object
networkPolicy:
properties:
metricsIngress:
items:
properties:
ipBlock:
properties:
cidr:
type: string
except:
items:
type: string
type: array
type: object
namespaceSelector:
properties:
matchExpressions:
items:
properties:
key:
type: string
operator:
type: string
values:
items:
type: string
type: array
type: object
type: array
matchLabels:
additionalProperties:
type: string
type: object
type: object
podSelector:
properties:
matchExpressions:
items:
properties:
key:
type: string
operator:
type: string
values:
items:
type: string
type: array
type: object
type: array
matchLabels:
additionalProperties:
type: string
type: object
type: object
type: object
minItems: 1
type: array
metricsPort:
maximum: 65535.0
minimum: 1.0
type: integer
required:
- metricsIngress
- metricsPort
type: object
serviceMetadata:
properties:
annotations:
Expand Down
52 changes: 52 additions & 0 deletions docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,58 @@ It is still honored: the NetworkPolicy is created when either key is `true`, so
`enable: true` in a base values file wins over `enabled: false` and must be removed to turn
the feature off. The same rule applies to `operatorConfiguration.dynamicConfig.enable`.

## Exposing SparkCluster Worker Metrics

Every `SparkCluster` gets a generated worker `NetworkPolicy` that only admits ingress from pods
carrying the cluster label or the driver-role label, so a Prometheus scraper is locked out by
default. Opening the worker web UI port (`8081` by default) is not a safe fix: Spark's built-in
`PrometheusServlet` metrics endpoint is served by the same embedded HTTP server as the web UI, so
admitting that port to any source would expose the whole UI, not just metrics.

Use the [Prometheus JMX Exporter](https://github.com/prometheus/jmx_exporter)
(`jmx_prometheus_javaagent`) to serve metrics on a dedicated HTTP port. See
[examples/cluster-with-jmx-exporter.yaml](../examples/cluster-with-jmx-exporter.yaml) for the full
`ConfigMap` and `SparkCluster` configuration. Replace the example's image with a custom Spark image
containing the exporter jar at `/opt/jmx_exporter/jmx_prometheus_javaagent.jar`; the stock
`apache/spark` image does not include it.

The example mounts the exporter rules and attaches the agent through `SPARK_DAEMON_JAVA_OPTS`.
The operator sets `SPARK_WORKER_OPTS` itself, so a value there would be overwritten. It also sets

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This describes the wrong failure mode. The operator does not replace a user-provided SPARK_WORKER_OPTS; it appends a second env entry with the same name (addNewEnv() in SparkClusterResourceSpec.java:340), and ClusterInitStep applies the StatefulSet with serverSideApply(), which rejects duplicate keys in the env associative list.

Checked against a live API server:

$ kubectl apply --server-side --dry-run=server -f dup-env.yaml
Error from server: failed to create typed patch object (default/dup-env-probe; apps/v1, Kind=StatefulSet):
  .spec.template.spec.containers[name="worker"].env: duplicate entries for key [name="SPARK_WORKER_OPTS"]

So a reader who takes "would be overwritten" at face value and sets it anyway does not get a silently ignored value. The whole worker StatefulSet apply is rejected, the cluster goes to SchedulingFailure and then Failed, and the status message is a stack trace that never mentions the duplicate variable.

The advice to use SPARK_DAEMON_JAVA_OPTS is right, only the reason is wrong. Something like: "SPARK_WORKER_OPTS must not be set there: the operator appends its own, and server-side apply rejects duplicate env names."

`spark.metrics.conf.worker.sink.jmx.class` to `org.apache.spark.metrics.sink.JmxSink` in `sparkConf` to
register Spark worker metrics as MBeans. Without this sink, the example exporter emits no Spark
worker metrics.
See also the [JmxSink configuration](spark_custom_resources.md) using `metrics.properties`.

Set `workerSpec.networkPolicy.metricsPort` to the exporter port and list the scraper's peers under
`metricsIngress`. The generated policy adds ingress on that port for those peers, alongside the
existing cluster/driver label allow-list. This mirrors
`operatorDeployment.networkPolicy.metricsIngress`
([above](#restricting-network-access-to-the-operator)) and accepts the same
[`NetworkPolicyPeer`](https://kubernetes.io/docs/reference/kubernetes-api/policy-resources/network-policy-v1/#NetworkPolicyPeer)
entries:

```yaml
spec:
workerSpec:
networkPolicy:
metricsPort: 9404
metricsIngress:
- namespaceSelector:
matchLabels:
kubernetes.io/metadata.name: "monitoring"
```

The `networkPolicy` block is optional. When present, the API server requires both fields and a
port between 1 and 65535, and rejects an empty `metricsIngress` list. Choose a dedicated exporter
port, not the worker web UI port; the operator does not verify this. Configuring this policy does
not start a metrics endpoint. The policy is generated when the cluster is created; changing scraper
peers requires recreating the `SparkCluster`.

Master pods do not get a `NetworkPolicy` today, so nothing needs to change for masters to be
scrapable — the same javaagent-plus-`ConfigMap` approach applied under `masterSpec` is enough to
expose master metrics, with no equivalent of `networkPolicy.metricsPort` /
`networkPolicy.metricsIngress` needed.

## Operator Health(Liveness) Probe with Sentinel Resource

Learning
Expand Down
71 changes: 71 additions & 0 deletions examples/cluster-with-jmx-exporter.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
apiVersion: v1
kind: ConfigMap
metadata:
name: jmx-exporter-config
data:
jmx-exporter-config.yaml: |
includeObjectNames: ["metrics:type=gauges,*"]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: the example asks the user to bake their own jmx_prometheus_javaagent jar but never names a version, and includeObjectNames is the newer key name. Older agent releases used whitelistObjectNames and reject unrecognized config keys outright. If someone picks an old jar, the agent fails to load this config, -javaagent initialization fails and the worker JVM aborts, so the pods CrashLoopBackOff while the SparkCluster still reports RunningHealthy (the operator sets that as soon as the applies succeed).

A minimum version next to the jar path comment would close this, e.g. "requires jmx_prometheus_javaagent 1.x; the includeObjectNames key is not accepted by older releases". I confirmed the key and the pattern work, but not which release introduced includeObjectNames, so please use whatever floor you are confident in.

rules:
- pattern: 'metrics<name=worker\.(\w+), type=gauges><>Value'
name: spark_worker_$1
type: GAUGE
---
apiVersion: spark.apache.org/v1
kind: SparkCluster
metadata:
name: cluster-with-jmx-exporter
spec:
runtimeVersions:
sparkVersion: "4.2.0"
clusterTolerations:
instanceConfig:
initWorkers: 1
minWorkers: 1
maxWorkers: 1
workerSpec:
networkPolicy:
metricsPort: 9404
# Allow scraper pods in "monitoring", in addition to the default cluster/driver peers.
metricsIngress:
- namespaceSelector:
matchLabels:
kubernetes.io/metadata.name: "monitoring"
statefulSetSpec:
template:
spec:
containers:
- name: worker
# Requires a custom image that has jmx_prometheus_javaagent pre-baked at this path
# (the stock apache/spark image does not include it).
env:
- name: SPARK_DAEMON_JAVA_OPTS
value: "-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=9404:/etc/metrics/jmx-exporter-config.yaml"
ports:
- name: jmx-metrics
containerPort: 9404
volumeMounts:
- name: jmx-exporter-config
mountPath: /etc/metrics
readOnly: true
volumes:
- name: jmx-exporter-config
configMap:
name: jmx-exporter-config
sparkConf:
spark.metrics.conf.worker.sink.jmx.class: "org.apache.spark.metrics.sink.JmxSink"
# Replace with a custom Spark image containing /opt/jmx_exporter/jmx_prometheus_javaagent.jar.
spark.kubernetes.container.image: "apache/spark:{{SPARK_VERSION}}-scala"
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.spark.k8s.operator.spec;

import java.util.List;

import com.fasterxml.jackson.annotation.JsonInclude;
import io.fabric8.generator.annotation.Max;
import io.fabric8.generator.annotation.Min;
import io.fabric8.generator.annotation.Required;
import io.fabric8.generator.annotation.Size;
import io.fabric8.kubernetes.api.model.networking.v1.NetworkPolicyPeer;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;

/**
* Network policy for the Spark workers.
*
* @since 1.1.0
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@JsonInclude(JsonInclude.Include.NON_NULL)
public class WorkerNetworkPolicySpec {
/**
* Port of a metrics endpoint (e.g. a JMX-to-Prometheus exporter agent) running in the worker
* container. This should be a dedicated exporter port, not the worker web UI port; the operator
* does not verify this. The caller is responsible for making the container listen on this port.
*/
@Required
@Min(1)
@Max(65535)
protected Integer metricsPort;
Comment thread
yalindogusahin marked this conversation as resolved.

/**
* Sources allowed to scrape {@link #metricsPort} on the worker pods, in addition to the sources
* the generated worker NetworkPolicy admits by default. At least one peer is required.
*/
@Required
@Size(min = 1)
protected List<NetworkPolicyPeer> metricsIngress;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

{metricsPort: 9404, metricsIngress: []} still passes the schema and silently does nothing, which is the same silent no-op we removed for partial blocks. In the Helm chart an empty list means "deny" around a metrics rule that always exists; here the block's only purpose is to add a rule, so an empty list has no meaning. @Size(min = 1) (it is in the same generator-annotations jar) emits minItems: 1, so the block is either effective or rejected at apply time. Please keep the Java isEmpty() guard regardless, see the comment there.

}
Original file line number Diff line number Diff line change
Expand Up @@ -45,4 +45,5 @@ public class WorkerSpec {
protected ServiceSpec serviceSpec;
protected ObjectMeta serviceMetadata;
protected HorizontalPodAutoscalerSpec horizontalPodAutoscalerSpec;
protected WorkerNetworkPolicySpec networkPolicy;
}
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,14 @@ public ReconcileProgress reconcile(
context.getClient().services().resource(masterService).forceConflicts().serverSideApply();
Service workerService = context.getWorkerServiceSpec();
context.getClient().services().resource(workerService).forceConflicts().serverSideApply();
NetworkPolicy workerNetworkPolicy = context.getWorkerNetworkPolicySpec();
context
.getClient()
.network()
.networkPolicies()
.resource(workerNetworkPolicy)
.forceConflicts()
.serverSideApply();
StatefulSet masterStatefulSet = context.getMasterStatefulSetSpec();
context
.getClient()
Expand All @@ -97,14 +105,6 @@ public ReconcileProgress reconcile(
.resource(workerStatefulSet)
.forceConflicts()
.serverSideApply();
NetworkPolicy workerNetworkPolicy = context.getWorkerNetworkPolicySpec();
context
.getClient()
.network()
.networkPolicies()
.resource(workerNetworkPolicy)
.forceConflicts()
.serverSideApply();
var horizontalPodAutoscaler = context.getHorizontalPodAutoscalerSpec();
if (horizontalPodAutoscaler.isPresent()) {
context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,47 @@ void suspendAfterMasterRequestedCompletesInitialization() {
verify(statefulSetApplicable, times(2)).serverSideApply();
}

@Test
void invalidWorkerNetworkPolicyDoesNotCreateStatefulSets() {
ClusterInitStep clusterInitStep = new ClusterInitStep();
SparkClusterContext mockContext = mock(SparkClusterContext.class);
SparkClusterStatusRecorder recorder = mock(SparkClusterStatusRecorder.class);
SparkCluster cluster = buildCluster();
KubernetesClient mockClient = mock(KubernetesClient.class, RETURNS_DEEP_STUBS);
when(mockContext.getResource()).thenReturn(cluster);
when(mockContext.getClient()).thenReturn(mockClient);
when(mockContext.getMasterServiceSpec()).thenReturn(service("cluster1-master-svc"));
when(mockContext.getWorkerServiceSpec()).thenReturn(service("cluster1-worker-svc"));
when(mockContext.getWorkerNetworkPolicySpec()).thenReturn(networkPolicy("cluster1-worker"));
@SuppressWarnings("unchecked")
ServerSideApplicable<Service> serviceApplicable = mock(ServerSideApplicable.class);
@SuppressWarnings("unchecked")
ServiceResource<Service> serviceResource = mock(ServiceResource.class);
when(serviceResource.forceConflicts()).thenReturn(serviceApplicable);
when(mockClient.services().resource(any(Service.class))).thenReturn(serviceResource);
@SuppressWarnings("unchecked")
ServerSideApplicable<NetworkPolicy> policyApplicable = mock(ServerSideApplicable.class);
@SuppressWarnings("unchecked")
Resource<NetworkPolicy> policyResource = mock(Resource.class);
when(policyResource.forceConflicts()).thenReturn(policyApplicable);
when(mockClient.network().networkPolicies().resource(any(NetworkPolicy.class)))
.thenReturn(policyResource);
when(policyApplicable.serverSideApply())
.thenThrow(new IllegalArgumentException("Invalid NetworkPolicyPeer"));

ReconcileProgress progress = clusterInitStep.reconcile(mockContext, recorder);

Assertions.assertEquals(ReconcileProgress.completeAndImmediateRequeue(), progress);
verify(mockClient, never()).apps();
verify(mockContext, never()).getMasterStatefulSetSpec();
verify(mockContext, never()).getWorkerStatefulSetSpec();
ArgumentCaptor<ClusterStatus> statusCaptor = ArgumentCaptor.forClass(ClusterStatus.class);
verify(recorder).persistStatus(any(), statusCaptor.capture());
Assertions.assertEquals(
ClusterStateSummary.SchedulingFailure,
statusCaptor.getValue().getCurrentState().getCurrentStateSummary());
}

@Test
void nonInitializingClusterProceeds() {
ClusterInitStep clusterInitStep = new ClusterInitStep();
Expand Down
Loading