4141import io .grpc .util .ForwardingSubchannel ;
4242import io .grpc .util .MultiChildLoadBalancer ;
4343import io .grpc .xds .internal .MetricReportUtils ;
44+ import io .grpc .xds .internal .MetricReportUtils .ParsedMetricName ;
4445import io .grpc .xds .orca .OrcaOobUtil ;
4546import io .grpc .xds .orca .OrcaOobUtil .OrcaOobReportListener ;
4647import io .grpc .xds .orca .OrcaPerRequestUtil ;
5051import java .util .HashSet ;
5152import java .util .List ;
5253import java .util .Objects ;
53- import java .util .OptionalDouble ;
5454import java .util .Random ;
5555import java .util .Set ;
5656import java .util .concurrent .ScheduledExecutorService ;
@@ -239,7 +239,7 @@ protected void updateOverallBalancingState() {
239239 private SubchannelPicker createReadyPicker (Collection <ChildLbState > activeList ) {
240240 WeightedRoundRobinPicker picker = new WeightedRoundRobinPicker (ImmutableList .copyOf (activeList ),
241241 config .enableOobLoadReport , config .errorUtilizationPenalty , sequence ,
242- config .metricNamesForComputingUtilization );
242+ config .parsedMetricNamesForComputingUtilization );
243243 updateWeight (picker );
244244 return picker ;
245245 }
@@ -329,15 +329,15 @@ public void addSubchannel(WrrSubchannel wrrSubchannel) {
329329 }
330330
331331 public OrcaReportListener getOrCreateOrcaListener (float errorUtilizationPenalty ,
332- ImmutableList <String > metricNamesForComputingUtilization ) {
332+ ImmutableList <ParsedMetricName > parsedMetricNamesForComputingUtilization ) {
333333 if (orcaReportListener != null
334334 && orcaReportListener .errorUtilizationPenalty == errorUtilizationPenalty
335- && orcaReportListener .metricNamesForComputingUtilization
336- .equals (metricNamesForComputingUtilization )) {
335+ && orcaReportListener .parsedMetricNamesForComputingUtilization
336+ .equals (parsedMetricNamesForComputingUtilization )) {
337337 return orcaReportListener ;
338338 }
339339 orcaReportListener =
340- new OrcaReportListener (errorUtilizationPenalty , metricNamesForComputingUtilization );
340+ new OrcaReportListener (errorUtilizationPenalty , parsedMetricNamesForComputingUtilization );
341341 return orcaReportListener ;
342342 }
343343
@@ -362,17 +362,17 @@ public void updateBalancingState(ConnectivityState newState, SubchannelPicker ne
362362
363363 final class OrcaReportListener implements OrcaPerRequestReportListener , OrcaOobReportListener {
364364 private final float errorUtilizationPenalty ;
365- private final ImmutableList <String > metricNamesForComputingUtilization ;
365+ private final ImmutableList <ParsedMetricName > parsedMetricNamesForComputingUtilization ;
366366
367367 OrcaReportListener (float errorUtilizationPenalty ,
368- ImmutableList <String > metricNamesForComputingUtilization ) {
368+ ImmutableList <ParsedMetricName > parsedMetricNamesForComputingUtilization ) {
369369 this .errorUtilizationPenalty = errorUtilizationPenalty ;
370- this .metricNamesForComputingUtilization = metricNamesForComputingUtilization ;
370+ this .parsedMetricNamesForComputingUtilization = parsedMetricNamesForComputingUtilization ;
371371 }
372372
373373 @ Override
374374 public void onLoadReport (MetricReport report ) {
375- double utilization = getUtilization (report , metricNamesForComputingUtilization );
375+ double utilization = getUtilization (report );
376376
377377 double newWeight = 0 ;
378378 if (utilization > 0 && report .getQps () > 0 ) {
@@ -398,10 +398,10 @@ public void onLoadReport(MetricReport report) {
398398 * if application utilization is > 0, it is returned. If neither are present, the CPU
399399 * utilization is returned.
400400 */
401- private double getUtilization (MetricReport report , ImmutableList < String > metricNames ) {
402- OptionalDouble customUtil = getCustomMetricUtilization (report , metricNames );
403- if (customUtil . isPresent () ) {
404- return customUtil . getAsDouble () ;
401+ private double getUtilization (MetricReport report ) {
402+ double customUtil = getCustomMetricUtilization (report );
403+ if (customUtil >= 0 ) {
404+ return customUtil ;
405405 }
406406 double appUtil = report .getApplicationUtilization ();
407407 if (appUtil > 0 ) {
@@ -412,20 +412,23 @@ private double getUtilization(MetricReport report, ImmutableList<String> metricN
412412
413413 /**
414414 * Returns the maximum utilization value among the specified metric names.
415- * Returns OptionalDouble.empty() if NONE of the specified metrics are present in the report,
415+ * Returns -1 if NONE of the specified metrics are present in the report,
416416 * or if all present metrics are NaN.
417- * Returns OptionalDouble.of(maxUtil) if at least one non-NaN metric is present.
418417 */
419- private OptionalDouble getCustomMetricUtilization (MetricReport report ,
420- ImmutableList <String > metricNames ) {
421- return metricNames .stream ()
422- .map (name -> MetricReportUtils .getMetric (report , name ))
423- .filter (OptionalDouble ::isPresent )
424- .mapToDouble (OptionalDouble ::getAsDouble )
425- .filter (d -> !Double .isNaN (d ) && d > 0 )
426- .max ();
418+ private double getCustomMetricUtilization (MetricReport report ) {
419+ double max = -1.0 ;
420+ for (int i = 0 ; i < parsedMetricNamesForComputingUtilization .size (); i ++) {
421+ double d = MetricReportUtils .getMetricValue (report ,
422+ parsedMetricNamesForComputingUtilization .get (i ));
423+ if (!Double .isNaN (d ) && d > 0 && d > max ) {
424+ max = d ;
425+ }
426+ }
427+ return max ;
427428 }
429+
428430 }
431+
429432 }
430433
431434 private final class UpdateWeightTask implements Runnable {
@@ -446,7 +449,7 @@ private void createAndApplyOrcaListeners() {
446449 if (config .enableOobLoadReport ) {
447450 OrcaOobUtil .setListener (weightedSubchannel ,
448451 wChild .getOrCreateOrcaListener (config .errorUtilizationPenalty ,
449- config .metricNamesForComputingUtilization ),
452+ config .parsedMetricNamesForComputingUtilization ),
450453 OrcaOobUtil .OrcaReportingConfig .newBuilder ()
451454 .setReportInterval (config .oobReportingPeriodNanos , TimeUnit .NANOSECONDS ).build ());
452455 } else {
@@ -516,7 +519,7 @@ static final class WeightedRoundRobinPicker extends SubchannelPicker {
516519
517520 WeightedRoundRobinPicker (List <ChildLbState > children , boolean enableOobLoadReport ,
518521 float errorUtilizationPenalty , AtomicInteger sequence ,
519- ImmutableList <String > metricNamesForComputingUtilization ) {
522+ ImmutableList <ParsedMetricName > parsedMetricNamesForComputingUtilization ) {
520523 checkNotNull (children , "children" );
521524 Preconditions .checkArgument (!children .isEmpty (), "empty child list" );
522525 this .children = children ;
@@ -526,7 +529,7 @@ static final class WeightedRoundRobinPicker extends SubchannelPicker {
526529 WeightedChildLbState wChild = (WeightedChildLbState ) child ;
527530 pickers .add (wChild .getCurrentPicker ());
528531 reportListeners .add (wChild .getOrCreateOrcaListener (errorUtilizationPenalty ,
529- metricNamesForComputingUtilization ));
532+ parsedMetricNamesForComputingUtilization ));
530533 }
531534 this .pickers = pickers ;
532535 this .reportListeners = reportListeners ;
@@ -768,6 +771,7 @@ static final class WeightedRoundRobinLoadBalancerConfig {
768771 final long weightUpdatePeriodNanos ;
769772 final float errorUtilizationPenalty ;
770773 final ImmutableList <String > metricNamesForComputingUtilization ;
774+ final ImmutableList <ParsedMetricName > parsedMetricNamesForComputingUtilization ;
771775
772776 public static Builder newBuilder () {
773777 return new Builder ();
@@ -784,6 +788,20 @@ private WeightedRoundRobinLoadBalancerConfig(long blackoutPeriodNanos,
784788 this .weightUpdatePeriodNanos = weightUpdatePeriodNanos ;
785789 this .errorUtilizationPenalty = errorUtilizationPenalty ;
786790 this .metricNamesForComputingUtilization = metricNamesForComputingUtilization ;
791+
792+ ImmutableList .Builder <ParsedMetricName > builder = ImmutableList .builder ();
793+ if (metricNamesForComputingUtilization != null ) {
794+ for (int i = 0 ; i < metricNamesForComputingUtilization .size (); i ++) {
795+ String metricName = metricNamesForComputingUtilization .get (i );
796+ ParsedMetricName parsed = MetricReportUtils .ParsedMetricName .parse (metricName );
797+ if (parsed .getMetricType () != MetricReportUtils .MetricType .INVALID ) {
798+ builder .add (parsed );
799+ } else {
800+ log .log (Level .FINE , "Invalid custom metric name configured and ignored: " + metricName );
801+ }
802+ }
803+ }
804+ this .parsedMetricNamesForComputingUtilization = builder .build ();
787805 }
788806
789807 @ Override
0 commit comments