Skip to content
Merged
Changes from 1 commit
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
46 changes: 34 additions & 12 deletions src/Orleans.Runtime/Placement/ActivationCountPlacementDirector.cs
Original file line number Diff line number Diff line change
Expand Up @@ -25,34 +25,44 @@ private class CachedLocalStat

// Track created activations on this silo between statistic intervals.
private readonly ConcurrentDictionary<SiloAddress, CachedLocalStat> _localCache = new();
private readonly SiloAddress _localAddress;
private readonly int _chooseHowMany;

public ActivationCountPlacementDirector(
ILocalSiloDetails localSiloDetails,
DeploymentLoadPublisher deploymentLoadPublisher,
IOptions<ActivationCountBasedPlacementOptions> options)
{
_localAddress = localSiloDetails.SiloAddress;
_ = localSiloDetails;
_chooseHowMany = options.Value.ChooseOutOf;
if (_chooseHowMany <= 0) throw new ArgumentException($"{nameof(ActivationCountBasedPlacementOptions)}.{nameof(ActivationCountBasedPlacementOptions.ChooseOutOf)} is {_chooseHowMany}. It must be greater than zero.");
deploymentLoadPublisher?.SubscribeToStatisticsChangeEvents(this);
}

private SiloAddress SelectSiloPowerOfK(SiloAddress[] silos)
{
var compatibleSilos = silos.ToSet();

// Exclude overloaded and non-compatible silos
var relevantSilos = new List<KeyValuePair<SiloAddress, CachedLocalStat>>();
var totalSilos = 0;
foreach (var kv in _localCache)
var totalSilos = _localCache.Count;
var compatibleSilosWithStats = 0;
var compatibleSilosWithoutStats = 0;
SiloAddress randomCompatibleSiloWithoutStats = default;
foreach (var silo in silos)
{
totalSilos++;
if (kv.Value.SiloStats.IsOverloaded) continue;
if (!compatibleSilos.Contains(kv.Key)) continue;
if (!_localCache.TryGetValue(silo, out var localSiloStat))
{
compatibleSilosWithoutStats++;
if (Random.Shared.Next(compatibleSilosWithoutStats) == 0)
{
randomCompatibleSiloWithoutStats = silo;
}

continue;
}

compatibleSilosWithStats++;
if (localSiloStat.SiloStats.IsOverloaded) continue;

relevantSilos.Add(kv);
relevantSilos.Add(new(silo, localSiloStat));
}

if (relevantSilos.Count > 0)
Expand Down Expand Up @@ -87,6 +97,18 @@ private SiloAddress SelectSiloPowerOfK(SiloAddress[] silos)
return minLoadedSilo.Key;
}

// If there are no stats for any compatible silos, fall back to random placement.
if (compatibleSilosWithStats == 0)
{
return silos[Random.Shared.Next(silos.Length)];
}
Comment thread
ReubenBond marked this conversation as resolved.
Outdated

// Some compatible silos might not have published statistics yet.
if (compatibleSilosWithoutStats > 0)
{
return randomCompatibleSiloWithoutStats;
}

// There are no compatible, non-overloaded silos.
var all = _localCache.ToList();
throw new SiloUnavailableException($"Unable to select a candidate from {all.Count} silos: {Utils.EnumerableToString(all, kvp => $"SiloAddress = {kvp.Key} -> {kvp.Value}")}");
Expand All @@ -104,10 +126,10 @@ private SiloAddress OnAddActivationInternal(PlacementTarget target, IPlacementCo
return placementHint;
}

// If the cache was not populated, just place locally
// If there are no statistics yet, fall back to random placement.
if (_localCache.IsEmpty)
{
return _localAddress;
return compatibleSilos[Random.Shared.Next(compatibleSilos.Length)];
}

return SelectSiloPowerOfK(compatibleSilos);
Expand Down
Loading