@@ -25,28 +25,28 @@ struct TShardMetaComp
2525 : Unit(Max(1UL , precisionBytes / blockSize))
2626 {}
2727
28- ui64 Score (const TShardBalancerBase::TShardMeta& m ) const
28+ ui64 Score (const TShardStats& stats ) const
2929 {
30- return Score (FreeSpace (m. Stats ));
30+ return Score (FreeSpace (stats ));
3131 }
3232
33- ui64 Score (ui64 s ) const
33+ ui64 Score (ui64 freeSpace ) const
3434 {
35- return static_cast <ui64>(round (s / static_cast <double >(Unit)));
35+ return static_cast <ui64>(round (freeSpace / static_cast <double >(Unit)));
3636 }
3737
3838 bool operator ()(
3939 const TShardBalancerBase::TShardMeta& lhs,
4040 const TShardBalancerBase::TShardMeta& rhs)
4141 {
42- return Score ( lhs) == Score ( rhs)
42+ return lhs. Score == rhs. Score
4343 ? lhs.ShardIdx < rhs.ShardIdx
44- : Score ( lhs) > Score ( rhs) ;
44+ : lhs. Score > rhs. Score ;
4545 }
4646
4747 bool operator ()(ui64 lhs, const TShardBalancerBase::TShardMeta& rhs)
4848 {
49- return Score (lhs) > Score ( rhs) ;
49+ return Score (lhs) > rhs. Score ;
5050 }
5151};
5252
@@ -64,19 +64,23 @@ TShardBalancerBase::TShardBalancerBase(
6464 , DesiredFreeSpaceReserve(desiredFreeSpaceReserve)
6565 , MinFreeSpaceReserve(minFreeSpaceReserve)
6666{
67+ // Before the first update shards are treated like empty with
68+ // infinite capacity. maxFileBlocks is used as infinity here.
69+ TShardMetaComp metaCmp (PrecisionBytes, BlockSize);
70+ Metas.reserve (shardIds.size ());
6771 for (ui32 i = 0 ; i < shardIds.size (); ++i) {
72+ TShardStats shardStats {
73+ .ShardId = shardIds[i],
74+ .TotalBlocksCount = maxFileBlocks,
75+ .UsedBlocksCount = 0 ,
76+ .UsedNodesCount = 0 ,
77+ .CurrentLoad = 0 ,
78+ .Suffer = 0 ,
79+ };
6880 Metas.emplace_back (
6981 i,
70- // Before the first update shards are treated like empty with
71- // infinite capacity. maxFileBlocks is used as infinity here.
72- TShardStats{
73- .ShardId = shardIds[i],
74- .TotalBlocksCount = maxFileBlocks,
75- .UsedBlocksCount = 0 ,
76- .UsedNodesCount = 0 ,
77- .CurrentLoad = 0 ,
78- .Suffer = 0 ,
79- });
82+ shardStats,
83+ metaCmp.Score (shardStats));
8084 }
8185}
8286
@@ -97,10 +101,11 @@ NProto::TError TShardBalancerBase::Update(
97101 << stats.size () << " != Metas.size() " << Metas.size ());
98102 }
99103
104+ const TShardMetaComp metaCmp (PrecisionBytes, BlockSize);
100105 for (ui32 i = 0 ; i < stats.size (); ++i) {
101- Metas[i] = TShardMeta (i, stats[i]);
106+ Metas[i] = TShardMeta (i, stats[i], metaCmp. Score (stats[i]) );
102107 }
103- Sort (Metas.begin (), Metas.end (), TShardMetaComp (PrecisionBytes, BlockSize) );
108+ Sort (Metas.begin (), Metas.end (), metaCmp );
104109 return {};
105110}
106111
@@ -246,6 +251,159 @@ NProto::TError TShardBalancerWeightedRandom::SelectShard(
246251
247252// //////////////////////////////////////////////////////////////////////////////
248253
254+ TShardBalancerWeightedDeterministic::TShardBalancerWeightedDeterministic (
255+ ui32 blockSize,
256+ ui64 precisionBytes,
257+ ui32 maxFileBlocks,
258+ ui64 desiredFreeSpaceReserve,
259+ ui64 minFreeSpaceReserve,
260+ TVector<TString> shardIds)
261+ : TShardBalancerBase(
262+ blockSize,
263+ precisionBytes,
264+ maxFileBlocks,
265+ desiredFreeSpaceReserve,
266+ minFreeSpaceReserve,
267+ std::move (shardIds))
268+ {
269+ // Update with zero stats.
270+ TVector<TShardStats> shardStats (Metas.size ());
271+ for (ui64 i = 0 ; i < shardStats.size (); ++i) {
272+ shardStats[i] = Metas[i].Stats ;
273+ }
274+ Update (shardStats, {}, {});
275+
276+ LastSelectedShard = Metas.size () - 1 ;
277+ CurrentScore = MaxScore;
278+ }
279+
280+ bool TShardBalancerWeightedDeterministic::CalcScore (
281+ const TVector<TShardStats>& stats)
282+ {
283+ ui64 minBlocksCount = Max<ui64>();
284+ for (const auto & stat: stats) {
285+ minBlocksCount = Min<ui64>(stat.UsedBlocksCount , minBlocksCount);
286+ }
287+
288+ const ui64 maxBlocksCount =
289+ Max<ui64>(stats[0 ].TotalBlocksCount / stats.size (), minBlocksCount);
290+ ui64 step = Max<ui64>(
291+ (maxBlocksCount - minBlocksCount) / ScoreLevelsCount,
292+ PrecisionBytes / BlockSize, 1UL );
293+
294+ bool scoreChanged = (Metas.size () != stats.size ());
295+ Metas.resize (stats.size ());
296+ for (ui64 i = 0 ; i < stats.size (); ++i) {
297+ const ui64 delta = maxBlocksCount > stats[i].UsedBlocksCount
298+ ? maxBlocksCount - stats[i].UsedBlocksCount
299+ : 0 ;
300+ const ui64 score = Min<ui64>(MaxScore, delta / step);
301+ scoreChanged = scoreChanged || (Metas[i].Score != score);
302+ Metas[i] = {static_cast <ui32>(i), stats[i], score};
303+ }
304+
305+ return scoreChanged;
306+ }
307+
308+ void TShardBalancerWeightedDeterministic::CalcNextShard ()
309+ {
310+ NextShard.SetSizes (Metas.size (), ScoreLevelsCount);
311+ NextShard.FillEvery (Metas.size ());
312+
313+ for (ui64 score = 0 ; score < ScoreLevelsCount; ++score) {
314+ ui64 firstQualifying = Metas.size ();
315+ for (ui64 shardIdx = 0 ; shardIdx < Metas.size (); ++shardIdx) {
316+ if (Metas[shardIdx].Score >= score) {
317+ firstQualifying = shardIdx;
318+ break ;
319+ }
320+ }
321+ if (firstQualifying == Metas.size ()) {
322+ continue ;
323+ }
324+
325+ ui64 nextQualifying = firstQualifying;
326+ for (ui64 shardIdx = Metas.size (); shardIdx-- > 0 ;) {
327+ NextShard[score][shardIdx] = nextQualifying;
328+ if (Metas[shardIdx].Score >= score) {
329+ nextQualifying = shardIdx;
330+ }
331+ }
332+ }
333+ }
334+
335+ NProto::TError TShardBalancerWeightedDeterministic::Update (
336+ const TVector<TShardStats>& stats,
337+ std::optional<ui64> desiredFreeSpaceReserve,
338+ std::optional<ui64> minFreeSpaceReserve)
339+ {
340+ Y_UNUSED (desiredFreeSpaceReserve);
341+ Y_UNUSED (minFreeSpaceReserve);
342+
343+ if (stats.size () != Metas.size ()) {
344+ return MakeError (E_ARGUMENT , TStringBuilder () << " stats.size() "
345+ << stats.size () << " != Metas.size() " << Metas.size ());
346+ }
347+
348+ if (!stats.empty () && CalcScore (stats)) {
349+ CalcNextShard ();
350+ }
351+
352+ return {};
353+ }
354+
355+ NProto::TError TShardBalancerWeightedDeterministic::SelectShard (
356+ ui64 fileSize,
357+ TString* shardId)
358+ {
359+ Y_UNUSED (fileSize);
360+
361+ if (Metas.empty ()) {
362+ return MakeError (
363+ E_ARGUMENT ,
364+ TStringBuilder () << " Metas.size() is zero" );
365+ }
366+
367+ // If jumping to the next shard crosses the right boundary,
368+ // increment CurrentScore.
369+ const ui64 nextShardIdx = NextShard[CurrentScore][LastSelectedShard];
370+ if (nextShardIdx < Metas.size () && nextShardIdx <= LastSelectedShard) {
371+ CurrentScore = (++CurrentScore) % ScoreLevelsCount;
372+ LastSelectedShard = Metas.size () - 1 ;
373+ }
374+
375+ // We are in an undefined part of the NextShard matrix.
376+ // That means that we should start from the initial state.
377+ if (NextShard[CurrentScore][LastSelectedShard] == Metas.size ()) {
378+ CurrentScore = 0 ;
379+ LastSelectedShard = Metas.size () - 1 ;
380+ }
381+
382+ LastSelectedShard = NextShard[CurrentScore][LastSelectedShard];
383+ *shardId = Metas[LastSelectedShard].Stats .ShardId ;
384+
385+ return {};
386+ }
387+
388+ TVector<TShardStats>
389+ TShardBalancerWeightedDeterministic::MakeOrderedShardList () const
390+ {
391+ TVector<TShardMeta> metas (Metas);
392+
393+ TShardMetaComp metaCmp (PrecisionBytes, BlockSize);
394+ Sort (metas.begin (), metas.end (), metaCmp);
395+
396+ TVector<TShardStats> shardStats;
397+ shardStats.reserve (metas.size ());
398+ for (const auto & meta: metas) {
399+ shardStats.emplace_back (meta.Stats );
400+ }
401+
402+ return shardStats;
403+ }
404+
405+ // //////////////////////////////////////////////////////////////////////////////
406+
249407IShardBalancerPtr CreateShardBalancer (
250408 NProto::EShardBalancerPolicy policy,
251409 ui32 blockSize,
@@ -280,6 +438,14 @@ IShardBalancerPtr CreateShardBalancer(
280438 desiredFreeSpaceReserve,
281439 minFreeSpaceReserve,
282440 std::move (shardIds));
441+ case NProto::SBP_WEIGHTED_DETERMINISTIC :
442+ return std::make_shared<TShardBalancerWeightedDeterministic>(
443+ blockSize,
444+ precisionBytes,
445+ maxFileBlocks,
446+ desiredFreeSpaceReserve,
447+ minFreeSpaceReserve,
448+ std::move (shardIds));
283449 default :
284450 Y_ABORT (
285451 " unsupported shard balancer policy: %d" ,
0 commit comments