Skip to content

Commit b2064a5

Browse files
authored
feat: notify Metrics when a controller's event processor starts (#3485)
Add Metrics.eventProcessingStarted(Controller), invoked in Controller whenever the event processor begins accepting events (both on inline start and after a deferred start, e.g. once leader election succeeds), so implementations can measure operator startup latency. Wire it through AggregatedMetrics so composed Metrics instances all receive the callback, and implement it in MicrometerMetrics/V2 as a per-controller gauge recording JVM uptime at the moment processing starts.
1 parent 958d1d6 commit b2064a5

6 files changed

Lines changed: 78 additions & 0 deletions

File tree

micrometer-support/src/main/java/io/javaoperatorsdk/operator/monitoring/micrometer/MicrometerMetrics.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.javaoperatorsdk.operator.monitoring.micrometer;
1717

18+
import java.lang.management.ManagementFactory;
1819
import java.util.*;
1920
import java.util.concurrent.ConcurrentHashMap;
2021
import java.util.concurrent.Executors;
@@ -54,6 +55,7 @@ public class MicrometerMetrics implements Metrics {
5455
private static final String RECONCILIATIONS_STARTED = RECONCILIATIONS + "started";
5556
private static final String RECONCILIATIONS_EXECUTIONS = PREFIX + RECONCILIATIONS + "executions.";
5657
private static final String RECONCILIATIONS_QUEUE_SIZE = PREFIX + RECONCILIATIONS + "queue.size.";
58+
private static final String PROCESSING_STARTED_LATENCY = PREFIX + "processing.started.latency.";
5759
private static final String NAME = "name";
5860
private static final String NAMESPACE = "namespace";
5961
private static final String GROUP = "group";
@@ -147,6 +149,16 @@ public void controllerRegistered(Controller<? extends HasMetadata> controller) {
147149
gauges.put(controllerQueueName, controllerQueueSize);
148150
}
149151

152+
@Override
153+
public void eventProcessingStarted(Controller<? extends HasMetadata> controller) {
154+
final var configuration = controller.getConfiguration();
155+
final var name = configuration.getName();
156+
final var tags = new ArrayList<Tag>(3);
157+
addGVKTags(GroupVersionKind.gvkFor(configuration.getResourceClass()), tags, false);
158+
registry.gauge(
159+
PROCESSING_STARTED_LATENCY + name, tags, ManagementFactory.getRuntimeMXBean().getUptime());
160+
}
161+
150162
@Override
151163
public <T> T timeControllerExecution(ControllerExecution<T> execution) {
152164
final var name = execution.controllerName();

micrometer-support/src/main/java/io/javaoperatorsdk/operator/monitoring/micrometer/MicrometerMetricsV2.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.javaoperatorsdk.operator.monitoring.micrometer;
1717

18+
import java.lang.management.ManagementFactory;
1819
import java.time.Duration;
1920
import java.util.*;
2021
import java.util.concurrent.ConcurrentHashMap;
@@ -59,6 +60,7 @@ public class MicrometerMetricsV2 implements Metrics {
5960
public static final String RECONCILIATIONS_EXECUTIONS_GAUGE = RECONCILIATIONS + "active";
6061
public static final String RECONCILIATIONS_QUEUE_SIZE_GAUGE = RECONCILIATIONS + "queue";
6162
public static final String NUMBER_OF_RESOURCE_GAUGE = "custom_resources";
63+
public static final String PROCESSING_STARTED_LATENCY_GAUGE = "processing.started.latency";
6264

6365
public static final String RECONCILIATION_EXECUTION_DURATION =
6466
RECONCILIATIONS + "execution.duration";
@@ -145,6 +147,15 @@ public void controllerRegistered(Controller<? extends HasMetadata> controller) {
145147
executionTimers.put(name, timer);
146148
}
147149

150+
@Override
151+
public void eventProcessingStarted(Controller<? extends HasMetadata> controller) {
152+
final var name = controller.getConfiguration().getName();
153+
final var tags = new ArrayList<Tag>();
154+
addControllerNameTag(name, tags);
155+
registry.gauge(
156+
PROCESSING_STARTED_LATENCY_GAUGE, tags, ManagementFactory.getRuntimeMXBean().getUptime());
157+
}
158+
148159
private String numberOfResourcesRefName(String name) {
149160
return NUMBER_OF_RESOURCE_GAUGE + name;
150161
}

operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/monitoring/AggregatedMetrics.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,11 @@ public void controllerRegistered(Controller<? extends HasMetadata> controller) {
7070
metricsList.forEach(metrics -> metrics.controllerRegistered(controller));
7171
}
7272

73+
@Override
74+
public void eventProcessingStarted(Controller<? extends HasMetadata> controller) {
75+
metricsList.forEach(metrics -> metrics.eventProcessingStarted(controller));
76+
}
77+
7378
@Override
7479
public void eventReceived(Event event, Map<String, Object> metadata) {
7580
metricsList.forEach(metrics -> metrics.eventReceived(event, metadata));

operator-framework-core/src/main/java/io/javaoperatorsdk/operator/api/monitoring/Metrics.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,17 @@ public interface Metrics {
4141
*/
4242
default void controllerRegistered(Controller<? extends HasMetadata> controller) {}
4343

44+
/**
45+
* Called when the event processor of the specified controller has started and begins processing
46+
* events. Until this point, events received from the event sources are deferred rather than
47+
* reconciled (see the {@code EventProcessor}); in leader-election setups this is only called once
48+
* leadership has been acquired. This callback can be used to measure how long the operator takes
49+
* to start accepting events.
50+
*
51+
* @param controller the controller whose event processor has started
52+
*/
53+
default void eventProcessingStarted(Controller<? extends HasMetadata> controller) {}
54+
4455
/**
4556
* Called when an event has been accepted by the SDK from an event source, which would potentially
4657
* trigger the Reconciler.

operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/Controller.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -383,6 +383,7 @@ public synchronized void start(boolean startEventProcessor) throws OperatorExcep
383383
eventSourceManager.start();
384384
if (startEventProcessor) {
385385
eventProcessor.start();
386+
metrics.eventProcessingStarted(this);
386387
}
387388
log.info("'{}' controller started", controllerName);
388389
} catch (MissingCRDException e) {
@@ -429,6 +430,7 @@ public synchronized void changeNamespaces(Set<String> namespaces) {
429430

430431
public synchronized void startEventProcessing() {
431432
eventProcessor.start();
433+
metrics.eventProcessingStarted(this);
432434
log.info("Started event processing for controller: {}", configuration.getName());
433435
}
434436

operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/ControllerTest.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import io.javaoperatorsdk.operator.api.config.ConfigurationService;
2929
import io.javaoperatorsdk.operator.api.config.MockControllerConfiguration;
3030
import io.javaoperatorsdk.operator.api.config.workflow.WorkflowSpec;
31+
import io.javaoperatorsdk.operator.api.monitoring.Metrics;
3132
import io.javaoperatorsdk.operator.api.reconciler.Cleaner;
3233
import io.javaoperatorsdk.operator.api.reconciler.DefaultContext;
3334
import io.javaoperatorsdk.operator.api.reconciler.DeleteControl;
@@ -68,6 +69,42 @@ void crdShouldNotBeCheckedForNativeResources() {
6869
verify(client, never()).apiextensions();
6970
}
7071

72+
@Test
73+
void notifiesMetricsWhenEventProcessorStarts() {
74+
final var client = MockKubernetesClient.client(Secret.class);
75+
final var metrics = mock(Metrics.class);
76+
final var configurationService =
77+
ConfigurationService.newOverriddenConfigurationService(
78+
new BaseConfigurationService(), o -> o.withMetrics(metrics));
79+
final var configuration =
80+
MockControllerConfiguration.forResource(Secret.class, configurationService);
81+
final var controller = new Controller<Secret>(reconciler, configuration, client);
82+
83+
// Inline start (event processing started synchronously, e.g. no leader election).
84+
controller.start(true);
85+
verify(metrics, times(1)).eventProcessingStarted(controller);
86+
87+
// Deferred start (event processing started later, e.g. after acquiring leadership).
88+
controller.startEventProcessing();
89+
verify(metrics, times(2)).eventProcessingStarted(controller);
90+
}
91+
92+
@Test
93+
void doesNotNotifyMetricsWhenEventProcessorNotStarted() {
94+
final var client = MockKubernetesClient.client(Secret.class);
95+
final var metrics = mock(Metrics.class);
96+
final var configurationService =
97+
ConfigurationService.newOverriddenConfigurationService(
98+
new BaseConfigurationService(), o -> o.withMetrics(metrics));
99+
final var configuration =
100+
MockControllerConfiguration.forResource(Secret.class, configurationService);
101+
final var controller = new Controller<Secret>(reconciler, configuration, client);
102+
103+
// e.g. leader election enabled: event sources start but event processing is deferred.
104+
controller.start(false);
105+
verify(metrics, never()).eventProcessingStarted(controller);
106+
}
107+
71108
@Test
72109
void crdShouldNotBeCheckedForCustomResourcesIfDisabled() {
73110
final var client = MockKubernetesClient.client(TestCustomResource.class);

0 commit comments

Comments
 (0)