Skip to content

Commit 1be4cd5

Browse files
author
Binbing Hou
committed
Add opt-in per-submission-mode breakdown for user metrics
1 parent 06f5733 commit 1be4cd5

5 files changed

Lines changed: 203 additions & 7 deletions

File tree

‎genie-web/src/integTest/java/com/netflix/genie/web/data/services/impl/jpa/JpaPersistenceServiceImplJobsIntegrationTest.java‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1067,15 +1067,27 @@ void canGetJobMetadata() throws GenieException {
10671067
@Test
10681068
@DatabaseSetup("persistence/jobs/search.xml")
10691069
void canGetUserResourceSummaries() {
1070-
final Map<String, UserResourcesSummary> summaries = this.service.getUserResourcesSummaries(
1070+
// API-submitted (api=true) jobs.
1071+
final Map<String, UserResourcesSummary> apiSummaries = this.service.getUserResourcesSummaries(
10711072
JobStatus.getActiveStatuses(),
10721073
true
10731074
);
1074-
Assertions.assertThat(summaries.keySet()).contains("tgianos");
1075-
final UserResourcesSummary userResourcesSummary = summaries.get("tgianos");
1076-
Assertions.assertThat(userResourcesSummary.getUser()).isEqualTo("tgianos");
1077-
Assertions.assertThat(userResourcesSummary.getRunningJobsCount()).isEqualTo(2L);
1078-
Assertions.assertThat(userResourcesSummary.getUsedMemory()).isEqualTo(4096L);
1075+
Assertions.assertThat(apiSummaries.keySet()).contains("tgianos");
1076+
final UserResourcesSummary apiSummary = apiSummaries.get("tgianos");
1077+
Assertions.assertThat(apiSummary.getUser()).isEqualTo("tgianos");
1078+
Assertions.assertThat(apiSummary.getRunningJobsCount()).isEqualTo(2L);
1079+
Assertions.assertThat(apiSummary.getUsedMemory()).isEqualTo(4096L);
1080+
1081+
// Agent/CLI-submitted (api=false) jobs are reported separately, with memory populated.
1082+
final Map<String, UserResourcesSummary> agentSummaries = this.service.getUserResourcesSummaries(
1083+
JobStatus.getActiveStatuses(),
1084+
false
1085+
);
1086+
Assertions.assertThat(agentSummaries.keySet()).contains("tgianos");
1087+
final UserResourcesSummary agentSummary = agentSummaries.get("tgianos");
1088+
Assertions.assertThat(agentSummary.getUser()).isEqualTo("tgianos");
1089+
Assertions.assertThat(agentSummary.getRunningJobsCount()).isEqualTo(2L);
1090+
Assertions.assertThat(agentSummary.getUsedMemory()).isEqualTo(4096L);
10791091
}
10801092

10811093
@Test

‎genie-web/src/main/java/com/netflix/genie/web/properties/UserMetricsProperties.java‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,4 +47,13 @@ public class UserMetricsProperties {
4747
private boolean enabled = true;
4848

4949
private long refreshInterval = 30_000;
50+
51+
/**
52+
* Whether per-user metrics are additionally published broken down by job submission mode (tagged
53+
* {@code submissionMode=api} for REST vs {@code submissionMode=agent} for agent/CLI), under the
54+
* {@code *.by-submission-mode.gauge} metric names. The existing
55+
* {@code genie.user.active-jobs.gauge}/{@code genie.user.active-memory.gauge} gauges are unaffected
56+
* regardless of this setting.
57+
*/
58+
private boolean splitBySubmissionMode;
5059
}

‎genie-web/src/main/java/com/netflix/genie/web/tasks/leader/UserMetricsTask.java‎

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@
3131
import io.micrometer.core.instrument.MeterRegistry;
3232
import lombok.extern.slf4j.Slf4j;
3333

34+
import java.util.Collections;
3435
import java.util.Map;
3536
import java.util.Set;
3637

@@ -46,12 +47,19 @@ public class UserMetricsTask extends LeaderTask {
4647
private static final String USER_ACTIVE_JOBS_METRIC_NAME = "genie.user.active-jobs.gauge";
4748
private static final String USER_ACTIVE_MEMORY_METRIC_NAME = "genie.user.active-memory.gauge";
4849
private static final String USER_ACTIVE_USERS_METRIC_NAME = "genie.user.active-users.gauge";
50+
private static final String USER_ACTIVE_JOBS_BY_SUBMISSION_MODE_METRIC_NAME =
51+
"genie.user.active-jobs.by-submission-mode.gauge";
52+
private static final String USER_ACTIVE_MEMORY_BY_SUBMISSION_MODE_METRIC_NAME =
53+
"genie.user.active-memory.by-submission-mode.gauge";
4954
private static final UserResourcesRecord USER_RECORD_PLACEHOLDER = new UserResourcesRecord("nobody");
5055
private final MeterRegistry registry;
5156
private final PersistenceService persistenceService;
5257
private final UserMetricsProperties userMetricsProperties;
5358

5459
private final Map<String, UserResourcesRecord> userResourcesRecordMap = Maps.newHashMap();
60+
// Records for the optional per-submission-mode resource breakdown: api (true/false) -> user -> record.
61+
private final Map<Boolean, Map<String, UserResourcesRecord>> userResourcesBySubmissionModeRecordMap =
62+
Maps.newHashMap();
5563
private final AtomicDouble activeUsersCount;
5664

5765
/**
@@ -153,6 +161,13 @@ public void run() {
153161
).update(jobs, memory);
154162
}
155163

164+
// Optionally publish a finer-grained breakdown tagged by submission mode (submissionMode=api REST vs
165+
// submissionMode=agent agent/CLI). Emitted under separate metric names so the gauges above are unchanged.
166+
// Disabled by default.
167+
if (this.userMetricsProperties.isSplitBySubmissionMode()) {
168+
this.publishResourcesBySubmissionMode(summaries);
169+
}
170+
156171
log.debug("Done publishing user metrics");
157172
}
158173

@@ -166,6 +181,7 @@ public void cleanup() {
166181

167182
// Reset all users
168183
this.userResourcesRecordMap.clear();
184+
this.userResourcesBySubmissionModeRecordMap.clear();
169185

170186
// Reset active users count
171187
this.activeUsersCount.set(Double.NaN);
@@ -189,6 +205,70 @@ private Number getUsersCount() {
189205
return activeUsersCount.get();
190206
}
191207

208+
private void publishResourcesBySubmissionMode(final Map<String, UserResourcesSummary> apiSummaries) {
209+
// Reuse the api=true summaries already fetched above; only the agent (api=false) slice is
210+
// an additional query.
211+
this.updateSubmissionModeRecords(true, apiSummaries);
212+
this.updateSubmissionModeRecords(
213+
false,
214+
this.persistenceService.getUserResourcesSummaries(JobStatus.getActiveStatuses(), false)
215+
);
216+
}
217+
218+
private void updateSubmissionModeRecords(final boolean api, final Map<String, UserResourcesSummary> summaries) {
219+
final Map<String, UserResourcesRecord> recordsForMode =
220+
this.userResourcesBySubmissionModeRecordMap.computeIfAbsent(api, mode -> Maps.newHashMap());
221+
222+
final String submissionModeTagValue =
223+
api ? MetricsConstants.TagValues.SUBMISSION_MODE_API : MetricsConstants.TagValues.SUBMISSION_MODE_AGENT;
224+
225+
// Drop users no longer active in this mode so their gauges fall back to NaN.
226+
recordsForMode.keySet().retainAll(summaries.keySet());
227+
228+
for (final UserResourcesSummary summary : summaries.values()) {
229+
final String user = summary.getUser();
230+
recordsForMode.computeIfAbsent(
231+
user,
232+
userName -> {
233+
// Gauge creation is idempotent so re-registration on reappearance is harmless.
234+
Gauge.builder(
235+
USER_ACTIVE_JOBS_BY_SUBMISSION_MODE_METRIC_NAME,
236+
() -> this.getSubmissionModeJobCount(api, userName)
237+
)
238+
.tags(
239+
MetricsConstants.TagKeys.USER, userName,
240+
MetricsConstants.TagKeys.SUBMISSION_MODE, submissionModeTagValue
241+
)
242+
.register(this.registry);
243+
Gauge.builder(
244+
USER_ACTIVE_MEMORY_BY_SUBMISSION_MODE_METRIC_NAME,
245+
() -> this.getSubmissionModeMemoryAmount(api, userName)
246+
)
247+
.tags(
248+
MetricsConstants.TagKeys.USER, userName,
249+
MetricsConstants.TagKeys.SUBMISSION_MODE, submissionModeTagValue
250+
)
251+
.register(this.registry);
252+
return new UserResourcesRecord(userName);
253+
}
254+
).update(summary.getRunningJobsCount(), summary.getUsedMemory());
255+
}
256+
}
257+
258+
private Number getSubmissionModeJobCount(final boolean api, final String user) {
259+
return this.userResourcesBySubmissionModeRecordMap
260+
.getOrDefault(api, Collections.emptyMap())
261+
.getOrDefault(user, USER_RECORD_PLACEHOLDER)
262+
.jobCount.get();
263+
}
264+
265+
private Number getSubmissionModeMemoryAmount(final boolean api, final String user) {
266+
return this.userResourcesBySubmissionModeRecordMap
267+
.getOrDefault(api, Collections.emptyMap())
268+
.getOrDefault(user, USER_RECORD_PLACEHOLDER)
269+
.memoryAmount.get();
270+
}
271+
192272
private static class UserResourcesRecord {
193273
private final String userName;
194274
private final AtomicDouble jobCount = new AtomicDouble(Double.NaN);

‎genie-web/src/main/java/com/netflix/genie/web/util/MetricsConstants.java‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,11 @@ public static final class TagKeys {
8181
*/
8282
public static final String USER = "user";
8383

84+
/**
85+
* Key to tag the job submission mode ({@code api} for REST API, {@code agent} for agent/CLI).
86+
*/
87+
public static final String SUBMISSION_MODE = "submissionMode";
88+
8489
/**
8590
* Key to tag the user concurrent job limit.
8691
*/
@@ -122,6 +127,16 @@ public static final class TagValues {
122127
*/
123128
public static final String FAILURE = "failure";
124129

130+
/**
131+
* Tag value for jobs submitted through the REST API (used with TagKeys.SUBMISSION_MODE).
132+
*/
133+
public static final String SUBMISSION_MODE_API = "api";
134+
135+
/**
136+
* Tag value for jobs submitted through the agent/CLI (used with TagKeys.SUBMISSION_MODE).
137+
*/
138+
public static final String SUBMISSION_MODE_AGENT = "agent";
139+
125140
/**
126141
* Utility class private constructor.
127142
*/

‎genie-web/src/test/groovy/com/netflix/genie/web/tasks/leader/UserMetricsTaskSpec.groovy‎

Lines changed: 81 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,9 +172,77 @@ class UserMetricsTaskSpec extends Specification {
172172
measureMemory("boo") == Double.NaN
173173
}
174174

175+
def "Run publishes a per-submission-mode breakdown when split-by-submission-mode is enabled"() {
176+
setup:
177+
Map<String, UserResourcesSummary> apiSummaries = [
178+
"foo": new UserResourcesSummary("foo", 10, 1024),
179+
"bar": new UserResourcesSummary("bar", 20, 2048)
180+
]
181+
Map<String, UserResourcesSummary> agentSummaries = [
182+
"foo": new UserResourcesSummary("foo", 3, 256),
183+
"baz": new UserResourcesSummary("baz", 5, 512)
184+
]
185+
186+
when:
187+
this.task = new UserMetricsTask(this.registry, this.dataServices, this.userMetricProperties)
188+
189+
then:
190+
1 * registry.gauge(_ as Meter.Id, _, _ as ToDoubleFunction) >> {
191+
args -> return captureGauge(args[0] as Meter.Id, args[1] as Object, args[2] as ToDoubleFunction<Object>)
192+
}
193+
194+
when:
195+
this.task.run()
196+
197+
then:
198+
_ * userMetricProperties.isSplitBySubmissionMode() >> true
199+
1 * persistenceService.getUserResourcesSummaries(JobStatus.activeStatuses, true) >> apiSummaries
200+
1 * persistenceService.getUserResourcesSummaries(JobStatus.activeStatuses, false) >> agentSummaries
201+
_ * registry.gauge(_ as Meter.Id, _, _ as ToDoubleFunction) >> {
202+
args -> return captureGauge(args[0] as Meter.Id, args[1] as Object, args[2] as ToDoubleFunction<Object>)
203+
}
204+
205+
// The existing api=true-only gauges are still published, unchanged.
206+
measureJobs("foo") == 10
207+
measureMemory("bar") == 2048
208+
209+
// New per-submission-mode gauges: api (REST) slice.
210+
measureJobsBySubmissionMode("foo", "api") == 10
211+
measureJobsBySubmissionMode("bar", "api") == 20
212+
measureMemoryBySubmissionMode("foo", "api") == 1024
213+
measureMemoryBySubmissionMode("bar", "api") == 2048
214+
215+
// New per-submission-mode gauges: agent slice.
216+
measureJobsBySubmissionMode("foo", "agent") == 3
217+
measureJobsBySubmissionMode("baz", "agent") == 5
218+
measureMemoryBySubmissionMode("foo", "agent") == 256
219+
measureMemoryBySubmissionMode("baz", "agent") == 512
220+
221+
when: "a user drops out of a mode on the next cycle"
222+
this.task.run()
223+
224+
then:
225+
_ * userMetricProperties.isSplitBySubmissionMode() >> true
226+
1 * persistenceService.getUserResourcesSummaries(JobStatus.activeStatuses, true) >> apiSummaries
227+
1 * persistenceService.getUserResourcesSummaries(JobStatus.activeStatuses, false) >>
228+
(["foo": new UserResourcesSummary("foo", 3, 256)] as Map<String, UserResourcesSummary>)
229+
_ * registry.gauge(_ as Meter.Id, _, _ as ToDoubleFunction) >> {
230+
args -> return captureGauge(args[0] as Meter.Id, args[1] as Object, args[2] as ToDoubleFunction<Object>)
231+
}
232+
233+
// baz no longer has agent jobs, so its by-submission-mode gauges report NaN.
234+
measureJobsBySubmissionMode("baz", "agent") == Double.NaN
235+
measureMemoryBySubmissionMode("baz", "agent") == Double.NaN
236+
measureJobsBySubmissionMode("foo", "agent") == 3
237+
measureMemoryBySubmissionMode("foo", "agent") == 256
238+
}
239+
175240
Gauge captureGauge(final Meter.Id id, final Object obj, final ToDoubleFunction<Object> f) {
176241
String userTagValue = id.getTag(MetricsConstants.TagKeys.USER)
177-
String gaugeKey = id.getName() + (userTagValue == null ? "" : ("-" + userTagValue))
242+
String submissionModeTagValue = id.getTag(MetricsConstants.TagKeys.SUBMISSION_MODE)
243+
String gaugeKey = id.getName() +
244+
(userTagValue == null ? "" : ("-" + userTagValue)) +
245+
(submissionModeTagValue == null ? "" : ("-" + submissionModeTagValue))
178246
this.gaugesFunctions.put(gaugeKey, { -> f.applyAsDouble(obj) })
179247
return Mock(Gauge)
180248
}
@@ -192,4 +260,16 @@ class UserMetricsTaskSpec extends Specification {
192260
String gaugeKey = UserMetricsTask.USER_ACTIVE_MEMORY_METRIC_NAME + "-" + user
193261
return gaugesFunctions.get(gaugeKey).call()
194262
}
263+
264+
double measureJobsBySubmissionMode(String user, String submissionMode) {
265+
String gaugeKey =
266+
UserMetricsTask.USER_ACTIVE_JOBS_BY_SUBMISSION_MODE_METRIC_NAME + "-" + user + "-" + submissionMode
267+
return gaugesFunctions.get(gaugeKey).call()
268+
}
269+
270+
double measureMemoryBySubmissionMode(String user, String submissionMode) {
271+
String gaugeKey =
272+
UserMetricsTask.USER_ACTIVE_MEMORY_BY_SUBMISSION_MODE_METRIC_NAME + "-" + user + "-" + submissionMode
273+
return gaugesFunctions.get(gaugeKey).call()
274+
}
195275
}

0 commit comments

Comments
 (0)