diff options
| author | Daniil Cherednik <[email protected]> | 2023-05-05 11:09:01 +0300 |
|---|---|---|
| committer | Daniil Cherednik <[email protected]> | 2023-05-05 11:09:01 +0300 |
| commit | b5a989b16cafa8a3b3bc076f1097a0eda6f48c06 (patch) | |
| tree | 4da744117a5aab37758921fa43b95a3068e5aec1 /library/cpp/actors/core/harmonizer.cpp | |
| parent | fc1cffcfa7f0497a1f97b384a24bcbf23362f3be (diff) | |
Ydb stable 23-1-2623.1.26
x-stable-origin-commit: 22184a7e157553d447f17a2dffc4ea2d32dfd74d
Diffstat (limited to 'library/cpp/actors/core/harmonizer.cpp')
| -rw-r--r-- | library/cpp/actors/core/harmonizer.cpp | 149 |
1 files changed, 129 insertions, 20 deletions
diff --git a/library/cpp/actors/core/harmonizer.cpp b/library/cpp/actors/core/harmonizer.cpp index e2fd0c5f24a..a1838ee2b09 100644 --- a/library/cpp/actors/core/harmonizer.cpp +++ b/library/cpp/actors/core/harmonizer.cpp @@ -19,8 +19,8 @@ constexpr bool CheckBinaryPower(ui64 value) { return !(value & (value - 1)); } +template <ui8 HistoryBufferSize = 8> struct TValueHistory { - static constexpr ui64 HistoryBufferSize = 8; static_assert(CheckBinaryPower(HistoryBufferSize)); double History[HistoryBufferSize] = {0.0}; @@ -31,32 +31,87 @@ struct TValueHistory { ui64 AccumulatedTs = 0; template <bool WithTail=false> - double GetAvgPartForLastSeconds(ui8 seconds) { - double sum = AccumulatedUs; + double Accumulate(auto op, auto comb, ui8 seconds) { + double acc = AccumulatedUs; size_t idx = HistoryIdx; ui8 leftSeconds = seconds; + if constexpr (!WithTail) { + idx--; + leftSeconds--; + if (idx >= HistoryBufferSize) { + idx = HistoryBufferSize - 1; + } + acc = History[idx]; + } do { idx--; leftSeconds--; if (idx >= HistoryBufferSize) { idx = HistoryBufferSize - 1; } - if (WithTail || leftSeconds) { - sum += History[idx]; + if constexpr (WithTail) { + acc = op(acc, History[idx]); + } else if (leftSeconds) { + acc = op(acc, History[idx]); } else { ui64 tsInSecond = Us2Ts(1'000'000.0); - sum += History[idx] * (tsInSecond - AccumulatedTs) / tsInSecond; + acc = op(acc, History[idx] * (tsInSecond - AccumulatedTs) / tsInSecond); } } while (leftSeconds); - double duration = 1'000'000.0 * seconds + (WithTail ? Ts2Us(AccumulatedTs): 0.0); - double avg = sum / duration; - return avg; + double duration = 1'000'000.0 * seconds; + if constexpr (WithTail) { + duration += Ts2Us(AccumulatedTs); + } + return comb(acc, duration); + } + + template <bool WithTail=false> + double GetAvgPartForLastSeconds(ui8 seconds) { + auto sum = [](double acc, double value) { + return acc + value; + }; + auto avg = [](double sum, double duration) { + return sum / duration; + }; + return Accumulate<WithTail>(sum, avg, seconds); } double GetAvgPart() { return GetAvgPartForLastSeconds<true>(HistoryBufferSize); } + double GetMaxForLastSeconds(ui8 seconds) { + auto max = [](const double& acc, const double& value) { + return Max(acc, value); + }; + auto fst = [](const double& value, const double&) { return value; }; + return Accumulate<false>(max, fst, seconds); + } + + double GetMax() { + return GetMaxForLastSeconds(HistoryBufferSize); + } + + i64 GetMaxInt() { + return static_cast<i64>(GetMax()); + } + + double GetMinForLastSeconds(ui8 seconds) { + auto min = [](const double& acc, const double& value) { + return Min(acc, value); + }; + auto fst = [](const double& value, const double&) { return value; }; + return Accumulate<false>(min, fst, seconds); + } + + double GetMin() { + return GetMinForLastSeconds(HistoryBufferSize); + } + + i64 GetMinInt() { + return static_cast<i64>(GetMin()); + } + void Register(ui64 ts, double valueUs) { if (ts < LastTs) { LastTs = ts; @@ -101,8 +156,8 @@ struct TValueHistory { }; struct TThreadInfo { - TValueHistory Consumed; - TValueHistory Booked; + TValueHistory<8> Consumed; + TValueHistory<8> Booked; }; struct TPoolInfo { @@ -116,6 +171,8 @@ struct TPoolInfo { NMonitoring::TDynamicCounters::TCounterPtr AvgPingCounterWithSmallWindow; ui32 MaxAvgPingUs = 0; ui64 LastUpdateTs = 0; + ui64 NotEnoughCpuExecutions = 0; + ui64 NewNotEnoughCpuExecutions = 0; TAtomic LastFlags = 0; // 0 - isNeedy; 1 - isStarved; 2 - isHoggish TAtomic IncreasingThreadsByNeedyState = 0; @@ -123,12 +180,20 @@ struct TPoolInfo { TAtomic DecreasingThreadsByHoggishState = 0; TAtomic PotentialMaxThreadCount = 0; + TValueHistory<16> Consumed; + TValueHistory<16> Booked; + + TAtomic MaxConsumedCpu = 0; + TAtomic MinConsumedCpu = 0; + TAtomic MaxBookedCpu = 0; + TAtomic MinBookedCpu = 0; + bool IsBeingStopped(i16 threadIdx); double GetBooked(i16 threadIdx); double GetlastSecondPoolBooked(i16 threadIdx); double GetConsumed(i16 threadIdx); double GetlastSecondPoolConsumed(i16 threadIdx); - void PullStats(ui64 ts); + TCpuConsumption PullStats(ui64 ts); i16 GetThreadCount(); void SetThreadCount(i16 threadCount); bool IsAvgPingGood(); @@ -167,15 +232,26 @@ double TPoolInfo::GetlastSecondPoolConsumed(i16 threadIdx) { } #define UNROLL_HISTORY(history) (history)[0], (history)[1], (history)[2], (history)[3], (history)[4], (history)[5], (history)[6], (history)[7] -void TPoolInfo::PullStats(ui64 ts) { +TCpuConsumption TPoolInfo::PullStats(ui64 ts) { + TCpuConsumption acc; for (i16 threadIdx = 0; threadIdx < MaxThreadCount; ++threadIdx) { TThreadInfo &threadInfo = ThreadInfo[threadIdx]; TCpuConsumption cpuConsumption = Pool->GetThreadCpuConsumption(threadIdx); + acc.Add(cpuConsumption); threadInfo.Consumed.Register(ts, cpuConsumption.ConsumedUs); LWPROBE(SavedValues, Pool->PoolId, Pool->GetName(), "consumed", UNROLL_HISTORY(threadInfo.Consumed.History)); threadInfo.Booked.Register(ts, cpuConsumption.BookedUs); LWPROBE(SavedValues, Pool->PoolId, Pool->GetName(), "booked", UNROLL_HISTORY(threadInfo.Booked.History)); } + Consumed.Register(ts, acc.ConsumedUs); + RelaxedStore(&MaxConsumedCpu, Consumed.GetMaxInt()); + RelaxedStore(&MinConsumedCpu, Consumed.GetMinInt()); + Booked.Register(ts, acc.BookedUs); + RelaxedStore(&MaxBookedCpu, Booked.GetMaxInt()); + RelaxedStore(&MinBookedCpu, Booked.GetMinInt()); + NewNotEnoughCpuExecutions = acc.NotEnoughCpuExecutions - NotEnoughCpuExecutions; + NotEnoughCpuExecutions = acc.NotEnoughCpuExecutions; + return acc; } #undef UNROLL_HISTORY @@ -206,6 +282,14 @@ private: std::vector<TPoolInfo> Pools; std::vector<ui16> PriorityOrder; + TValueHistory<16> Consumed; + TValueHistory<16> Booked; + + TAtomic MaxConsumedCpu = 0; + TAtomic MinConsumedCpu = 0; + TAtomic MaxBookedCpu = 0; + TAtomic MinBookedCpu = 0; + void PullStats(ui64 ts); void HarmonizeImpl(ui64 ts); void CalculatePriorityOrder(); @@ -217,7 +301,8 @@ public: void DeclareEmergency(ui64 ts) override; void AddPool(IExecutorPool* pool, TSelfPingInfo *pingInfo) override; void Enable(bool enable) override; - TPoolHarmonizedStats GetPoolStats(i16 poolId) const override; + TPoolHarmonizerStats GetPoolStats(i16 poolId) const override; + THarmonizerStats GetStats() const override; }; THarmonizer::THarmonizer(ui64 ts) { @@ -232,13 +317,21 @@ double THarmonizer::Rescale(double value) const { } void THarmonizer::PullStats(ui64 ts) { + TCpuConsumption acc; for (TPoolInfo &pool : Pools) { - pool.PullStats(ts); + TCpuConsumption consumption = pool.PullStats(ts); + acc.Add(consumption); } + Consumed.Register(ts, acc.ConsumedUs); + RelaxedStore(&MaxConsumedCpu, Consumed.GetMaxInt()); + RelaxedStore(&MinConsumedCpu, Consumed.GetMinInt()); + Booked.Register(ts, acc.BookedUs); + RelaxedStore(&MaxBookedCpu, Booked.GetMaxInt()); + RelaxedStore(&MinBookedCpu, Booked.GetMinInt()); } Y_FORCE_INLINE bool IsStarved(double consumed, double booked) { - return consumed < booked * 0.7; + return Max(consumed, booked) > 0.1 && consumed < booked * 0.7; } Y_FORCE_INLINE bool IsHoggish(double booked, ui16 currentThreadCount) { @@ -273,7 +366,7 @@ void THarmonizer::HarmonizeImpl(ui64 ts) { isStarvedPresent = true; } ui32 currentThreadCount = pool.GetThreadCount(); - bool isNeedy = pool.IsAvgPingGood() && poolBooked >= currentThreadCount; + bool isNeedy = (pool.IsAvgPingGood() || pool.NewNotEnoughCpuExecutions) && poolBooked >= currentThreadCount; if (pool.AvgPingCounter) { if (pool.LastUpdateTs + Us2Ts(3'000'000ull) > ts) { isNeedy = false; @@ -304,6 +397,9 @@ void THarmonizer::HarmonizeImpl(ui64 ts) { AtomicSet(pool.PotentialMaxThreadCount, Min(pool.MaxThreadCount, budgetInt)); } double overbooked = consumed - booked; + if (overbooked < 0) { + isStarvedPresent = false; + } if (isStarvedPresent) { // last_starved_at_consumed_value = сумма по всем пулам consumed; // TODO(cthulhu): использовать как лимит планвно устремлять этот лимит к total, @@ -319,7 +415,7 @@ void THarmonizer::HarmonizeImpl(ui64 ts) { TPoolInfo &pool = Pools[poolIdx]; i64 threadCount = pool.GetThreadCount(); while (threadCount > pool.DefaultThreadCount) { - pool.SetThreadCount(threadCount - 1); + pool.SetThreadCount(--threadCount); AtomicIncrement(pool.DecreasingThreadsByStarvedState); overbooked--; LWPROBE(HarmonizeOperation, poolIdx, pool.Pool->GetName(), "decrease", threadCount - 1, pool.DefaultThreadCount, pool.MaxThreadCount); @@ -425,13 +521,17 @@ IHarmonizer* MakeHarmonizer(ui64 ts) { return new THarmonizer(ts); } -TPoolHarmonizedStats THarmonizer::GetPoolStats(i16 poolId) const { +TPoolHarmonizerStats THarmonizer::GetPoolStats(i16 poolId) const { const TPoolInfo &pool = Pools[poolId]; ui64 flags = RelaxedLoad(&pool.LastFlags); - return TPoolHarmonizedStats { + return TPoolHarmonizerStats{ .IncreasingThreadsByNeedyState = static_cast<ui64>(RelaxedLoad(&pool.IncreasingThreadsByNeedyState)), .DecreasingThreadsByStarvedState = static_cast<ui64>(RelaxedLoad(&pool.DecreasingThreadsByStarvedState)), .DecreasingThreadsByHoggishState = static_cast<ui64>(RelaxedLoad(&pool.DecreasingThreadsByHoggishState)), + .MaxConsumedCpu = static_cast<i64>(RelaxedLoad(&pool.MaxConsumedCpu)), + .MinConsumedCpu = static_cast<i64>(RelaxedLoad(&pool.MinConsumedCpu)), + .MaxBookedCpu = static_cast<i64>(RelaxedLoad(&pool.MaxBookedCpu)), + .MinBookedCpu = static_cast<i64>(RelaxedLoad(&pool.MinBookedCpu)), .PotentialMaxThreadCount = static_cast<i16>(RelaxedLoad(&pool.PotentialMaxThreadCount)), .IsNeedy = static_cast<bool>(flags & 1), .IsStarved = static_cast<bool>(flags & 2), @@ -439,4 +539,13 @@ TPoolHarmonizedStats THarmonizer::GetPoolStats(i16 poolId) const { }; } +THarmonizerStats THarmonizer::GetStats() const { + return THarmonizerStats{ + .MaxConsumedCpu = static_cast<i64>(RelaxedLoad(&MaxConsumedCpu)), + .MinConsumedCpu = static_cast<i64>(RelaxedLoad(&MinConsumedCpu)), + .MaxBookedCpu = static_cast<i64>(RelaxedLoad(&MaxBookedCpu)), + .MinBookedCpu = static_cast<i64>(RelaxedLoad(&MinBookedCpu)), + }; +} + } |
