From b6370ce25c33539c1ea6dfd9744d8072dc3f92ae Mon Sep 17 00:00:00 2001 From: Neil Date: Fri, 3 Apr 2026 11:34:02 +0100 Subject: [PATCH 1/2] chore: implement leastlatency LB with Peak EWMA Other changes - Hook up subset size settings - Increase PeakEwma decay - AFE latency is very noisy so we need a long duration to discern which AFEs genuinely perform better Change-Id: I138501487d4dec53aa80e18785023bcbfaa06807 --- .../v2/internal/session/DynamicPicker.java | 36 ++++++----- .../internal/session/LeastInFlightPicker.java | 37 +++++++---- .../internal/session/LeastLatencyPicker.java | 64 +++++++++++++++++++ .../data/v2/internal/session/SessionList.java | 4 +- .../v2/internal/session/SessionPoolImpl.java | 2 +- .../v2/internal/session/SimplePicker.java | 5 +- 6 files changed, 113 insertions(+), 35 deletions(-) create mode 100644 google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/DynamicPicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/DynamicPicker.java index 664f8e64657a..1d932a865b6c 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/DynamicPicker.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/DynamicPicker.java @@ -31,13 +31,12 @@ class DynamicPicker extends Picker { private final SessionList sessions; private volatile Picker delegate; - private LoadBalancingOptions.LoadBalancingStrategyCase currentStrategy; + private volatile LoadBalancingOptions currentOptions; - public DynamicPicker( - SessionList sessions, LoadBalancingOptions.LoadBalancingStrategyCase initialStrategy) { + public DynamicPicker(SessionList sessions, LoadBalancingOptions initialOptions) { this.sessions = sessions; - this.currentStrategy = initialStrategy; - this.delegate = createPicker(initialStrategy); + this.currentOptions = initialOptions; + this.delegate = createPicker(initialOptions); } @Override @@ -46,25 +45,28 @@ public Optional pickSession() { } public void updateConfig(SessionClientConfiguration.SessionPoolConfiguration config) { - LoadBalancingOptions.LoadBalancingStrategyCase newStrategy = - config.getLoadBalancingOptions().getLoadBalancingStrategyCase(); - if (newStrategy != currentStrategy) { - delegate = createPicker(newStrategy); - currentStrategy = newStrategy; + LoadBalancingOptions newOptions = config.getLoadBalancingOptions(); + if (!newOptions.equals(currentOptions)) { + delegate = createPicker(newOptions); + currentOptions = newOptions; } } - private Picker createPicker(LoadBalancingOptions.LoadBalancingStrategyCase strategy) { - switch (strategy) { + private Picker createPicker(LoadBalancingOptions options) { + switch (options.getLoadBalancingStrategyCase()) { case RANDOM: - return new SimplePicker(sessions); + return new SimplePicker(sessions, options.getRandom()); case LEAST_IN_FLIGHT: - return new LeastInFlightPicker(sessions); + return new LeastInFlightPicker(sessions, options.getLeastInFlight()); + case PEAK_EWMA: + return new LeastLatencyPicker(sessions, options.getPeakEwma()); default: LOGGER.log( - Level.FINE, "got load balancing strategy {0} which was not implemented", strategy); - // TODO: implement PeakEwma - return new LeastInFlightPicker(sessions); + Level.FINE, + "got load balancing strategy {0} which was not implemented", + options.getLoadBalancingStrategyCase()); + return new LeastInFlightPicker( + sessions, LoadBalancingOptions.LeastInFlight.getDefaultInstance()); } } } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java index 8a3e195b1fdb..86ec52a9ff91 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java @@ -16,40 +16,49 @@ package com.google.cloud.bigtable.data.v2.internal.session; +import com.google.bigtable.v2.LoadBalancingOptions; import com.google.cloud.bigtable.data.v2.internal.session.SessionList.AfeHandle; import com.google.cloud.bigtable.data.v2.internal.session.SessionList.SessionHandle; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.concurrent.ThreadLocalRandom; -/** Pick the AFE with the fewest in-flight requests. Experimental for now. */ +/** Pick the AFE with the fewest in-flight requests. */ class LeastInFlightPicker extends Picker { private final SessionList sessionList; + private final LoadBalancingOptions.LeastInFlight options; - public LeastInFlightPicker(SessionList sessionList) { + public LeastInFlightPicker(SessionList sessionList, LoadBalancingOptions.LeastInFlight options) { this.sessionList = sessionList; + this.options = options; } @Override Optional pickSession() { List readyAfes = sessionList.getAfesWithReadySessions(); - int size = readyAfes.size(); - - if (size == 0) { + if (readyAfes.isEmpty()) { return Optional.empty(); } - ThreadLocalRandom random = ThreadLocalRandom.current(); - AfeHandle selected = readyAfes.get(random.nextInt(size)); + ThreadLocalRandom rng = ThreadLocalRandom.current(); + List candidates = new ArrayList<>(readyAfes); + int bestCost = Integer.MAX_VALUE; + AfeHandle bestAfe = null; + long iterations = Math.min(options.getRandomSubsetSize(), readyAfes.size()); - // If we have options, pick a second candidate and keep the better one - if (size > 1) { - AfeHandle candidate2 = readyAfes.get(random.nextInt(size)); - if (candidate2.getNumOutstanding() < selected.getNumOutstanding()) { - selected = candidate2; + // Partial Fisher-Yates shuffle. + for (int i = 0; i < iterations; i++) { + int randomIndex = i + rng.nextInt(candidates.size() - i); + AfeHandle picked = candidates.get(randomIndex); + if (picked.getNumOutstanding() < bestCost) { + bestCost = picked.getNumOutstanding(); + bestAfe = picked; } + // Move candidate to the `i`th entry so that it's not picked again. + Collections.swap(candidates, i, randomIndex); } - - return sessionList.checkoutSession(selected); + return sessionList.checkoutSession(bestAfe); } } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java new file mode 100644 index 000000000000..8cbd2d265bc1 --- /dev/null +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java @@ -0,0 +1,64 @@ +/* + * Copyright 2026 Google LLC + * + * Licensed 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 + * + * https://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 com.google.cloud.bigtable.data.v2.internal.session; + +import com.google.bigtable.v2.LoadBalancingOptions; +import com.google.cloud.bigtable.data.v2.internal.session.SessionList.AfeHandle; +import com.google.cloud.bigtable.data.v2.internal.session.SessionList.SessionHandle; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.ThreadLocalRandom; + +/** Pick the AFE with the least latency. Experimental for now. */ +class LeastLatencyPicker extends Picker { + private final SessionList sessionList; + private final LoadBalancingOptions.PeakEwma options; + + public LeastLatencyPicker(SessionList sessionList, LoadBalancingOptions.PeakEwma options) { + this.sessionList = sessionList; + this.options = options; + } + + @Override + Optional pickSession() { + List readyAfes = sessionList.getAfesWithReadySessions(); + if (readyAfes.isEmpty()) { + return Optional.empty(); + } + + ThreadLocalRandom rng = ThreadLocalRandom.current(); + List candidates = new ArrayList<>(readyAfes); + double bestCost = Double.MAX_VALUE; + AfeHandle bestAfe = null; + long iterations = Math.min(options.getRandomSubsetSize(), readyAfes.size()); + + // Partial Fisher-Yates shuffle. + for (int i = 0; i < iterations; i++) { + int randomIndex = i + rng.nextInt(candidates.size() - i); + AfeHandle picked = candidates.get(randomIndex); + if (picked.getE2eCost() < bestCost) { + bestCost = picked.getE2eCost(); + bestAfe = picked; + } + // Move candidate to the `i`th entry so that it's not picked again. + Collections.swap(candidates, i, randomIndex); + } + return sessionList.checkoutSession(bestAfe); + } +} diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionList.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionList.java index 0cc06e258ad4..c5efeebffd7d 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionList.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionList.java @@ -413,8 +413,8 @@ int getNumOutstanding() { } static class PeakEwma { - // Use the last 100ms as a look back window - private final double decayNs = TimeUnit.MILLISECONDS.toNanos(100); + // Use the last 10s as a look back window + private final double decayNs = TimeUnit.SECONDS.toNanos(10); private long timestamp = System.nanoTime(); private double cost; diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionPoolImpl.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionPoolImpl.java index 9a69907965b8..98e9cd909d1d 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionPoolImpl.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SessionPoolImpl.java @@ -192,7 +192,7 @@ public SessionPoolImpl( .getSessionConfiguration() .getSessionPoolConfiguration() .getLoadBalancingOptions(); - picker = new DynamicPicker(sessions, lbOptions.getLoadBalancingStrategyCase()); + picker = new DynamicPicker(sessions, lbOptions); poolSizer = new PoolSizer( sessions.getStats(), diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SimplePicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SimplePicker.java index d80eb252291d..5a0acfb83bcc 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SimplePicker.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/SimplePicker.java @@ -16,6 +16,7 @@ package com.google.cloud.bigtable.data.v2.internal.session; +import com.google.bigtable.v2.LoadBalancingOptions; import com.google.cloud.bigtable.data.v2.internal.session.SessionList.AfeHandle; import com.google.cloud.bigtable.data.v2.internal.session.SessionList.SessionHandle; import java.util.List; @@ -24,10 +25,12 @@ class SimplePicker extends Picker { private final SessionList sessionList; + private final LoadBalancingOptions.Random options; private final Random random = new Random(); - public SimplePicker(SessionList sessionList) { + public SimplePicker(SessionList sessionList, LoadBalancingOptions.Random options) { this.sessionList = sessionList; + this.options = options; } @Override From 3633c9cb414579adba22277759db465c8dc545e7 Mon Sep 17 00:00:00 2001 From: Igor Bernstein Date: Thu, 9 Apr 2026 09:39:04 -0400 Subject: [PATCH 2/2] fix: handling default handling where subsetSize = 0 means the entire pool Change-Id: Ifb6c293e355d337d457d3b76efe4e21924f80f81 --- .../data/v2/internal/session/LeastInFlightPicker.java | 5 ++++- .../data/v2/internal/session/LeastLatencyPicker.java | 6 +++++- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java index 86ec52a9ff91..fe6cd7becabd 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastInFlightPicker.java @@ -46,7 +46,10 @@ Optional pickSession() { List candidates = new ArrayList<>(readyAfes); int bestCost = Integer.MAX_VALUE; AfeHandle bestAfe = null; - long iterations = Math.min(options.getRandomSubsetSize(), readyAfes.size()); + long iterations = readyAfes.size(); + if (options.getRandomSubsetSize() > 0) { + iterations = Math.min(options.getRandomSubsetSize(), iterations); + } // Partial Fisher-Yates shuffle. for (int i = 0; i < iterations; i++) { diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java index 8cbd2d265bc1..c4e264b1fa3b 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/session/LeastLatencyPicker.java @@ -46,7 +46,11 @@ Optional pickSession() { List candidates = new ArrayList<>(readyAfes); double bestCost = Double.MAX_VALUE; AfeHandle bestAfe = null; - long iterations = Math.min(options.getRandomSubsetSize(), readyAfes.size()); + long iterations = readyAfes.size(); + + if (options.getRandomSubsetSize() > 0) { + iterations = Math.min(options.getRandomSubsetSize(), iterations); + } // Partial Fisher-Yates shuffle. for (int i = 0; i < iterations; i++) {