Skip to content

Commit 1db0efc

Browse files
committed
fix(messagequeue): remainder-aware fair share and orphan sweep
## Summary ### Why? The MySQL queue subscriber capped every node at ceil(P/N) partitions, computed independently per node. The caps sum to more than P, so uneven splits can settle into stable starvation: with 12 partitions and 5 subscribers every cap is 3 (sum 15), and 3/3/3/3/0 is a legal end state where every owner sits "at cap" and nothing ever obliges anyone to shed for the empty subscriber. Separately, nothing guaranteed an unleased partition would eventually be picked up if the share arithmetic misfired (divergent heartbeat views during the staleness window, or a subscriber that heartbeats without ever acquiring), and every shedding rebalance tick logged a spurious `lease renewal failed: ErrLeaseExpired` because rebalance sorted the shared slice in place and renewal then ran over the just-released tail. ### What? - `fairShareCap` is now remainder-aware: each subscriber ranks itself in the sorted `ActiveSubscribers` list; the first `P mod N` ranks get `floor(P/N)+1` and the rest `floor(P/N)`, so per-rank caps sum to exactly P. A starved subscriber or an unclaimed partition now implies a peer over/under its cap that rebalance and discovery resolve — neither is a stable state. The signature, the `0 = unlimited` contract, and the minimum-1 floor are unchanged; a subscriber missing from its own active view falls back to `ceil` over N+1 contenders (never unlimited, never zero). - Orphan sweep: every `2 x LeaseDurationMs` the discovery tick runs one uncapped acquisition pass. `TryAcquireLease` cannot steal a valid lease, so the sweep is a no-op in a healthy group, but a partition left unleased for any reason is picked up by whichever subscriber sweeps first; the next rebalance sheds any over-cap grab once a peer has capacity. - `rebalance` sheds from a sorted copy instead of reordering the caller's slice, and returns the released partitions; the lease tick renews only the kept set, eliminating the spurious `ErrLeaseExpired` error log on every shedding tick. ## Test Plan - ✅ `make test` — new `TestSubscriber_FairShareCap` covers the 12/5 starvation case (caps 3,3,2,2,2), P<N flooring, the missing-heartbeat fallback, and a property subtest asserting per-rank caps sum to exactly P for all N≤6, P≤13; `TestSubscriber_Rebalance*` cover the released-tail contract and that the input slice is not reordered. - ✅ `bazel test //test/integration/extension/messagequeue/...` — full Docker suite, including new `TestRebalance_NoStarvation_UnevenSplit` (12 partitions / 5 subscribers converge with every subscriber owning 2–3, total 12) and `TestRebalance_OrphanSweep` (phantom heartbeats cap the real subscriber at 1 of 3 partitions; all 3 messages are still delivered and acked via the sweep).
1 parent 503ae22 commit 1db0efc

3 files changed

Lines changed: 447 additions & 23 deletions

File tree

‎platform/extension/messagequeue/mysql/subscriber.go‎

Lines changed: 103 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -485,6 +485,13 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
485485
leaseTicker := time.NewTicker(time.Duration(cfg.LeaseRenewalIntervalMs) * time.Millisecond)
486486
defer leaseTicker.Stop()
487487

488+
// Orphan sweep pacing: an uncapped acquisition pass runs every two lease
489+
// durations. Two lease durations is past every transient window in the
490+
// protocol — an expiring lease, a crashed peer's heartbeat going stale —
491+
// so anything still unleased at sweep time is genuinely unclaimed.
492+
orphanSweepInterval := 2 * time.Duration(cfg.LeaseDurationMs) * time.Millisecond
493+
lastOrphanSweep := time.Now()
494+
488495
// Send initial heartbeat so this subscriber is immediately visible to
489496
// ActiveSubscribers. Without this, other subscribers compute incorrect
490497
// fair shares until the first leaseTicker fires.
@@ -533,10 +540,26 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
533540

534541
// Rebalance, renew, and heartbeat are independent operations.
535542
// Each can fail without affecting the others — the next tick retries.
536-
if err := s.rebalance(ctx, sub, leasedPartitions); err != nil {
543+
// Renewal covers only the partitions kept after shedding; renewing
544+
// a just-released lease would spuriously fail with ErrLeaseExpired.
545+
released, err := s.rebalance(ctx, sub, leasedPartitions)
546+
if err != nil {
537547
s.logger.Errorw("rebalance failed", append(logFields, "error", err)...)
538548
}
539-
if err := s.renewLeases(ctx, sub, leasedPartitions); err != nil {
549+
kept := leasedPartitions
550+
if len(released) > 0 {
551+
releasedSet := make(map[string]struct{}, len(released))
552+
for _, pk := range released {
553+
releasedSet[pk] = struct{}{}
554+
}
555+
kept = make([]string, 0, len(leasedPartitions))
556+
for _, pk := range leasedPartitions {
557+
if _, ok := releasedSet[pk]; !ok {
558+
kept = append(kept, pk)
559+
}
560+
}
561+
}
562+
if err := s.renewLeases(ctx, sub, kept); err != nil {
540563
s.logger.Errorw("lease renewal failed", append(logFields, "error", err)...)
541564
}
542565
if err := s.sendHeartbeat(ctx, sub); err != nil {
@@ -545,7 +568,18 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
545568
s.emitSignal(SignalPartitionUpdate)
546569

547570
case <-discoveryTicker.C:
548-
if err := s.discoverAndReconcileWorkers(ctx, sub); err != nil {
571+
// Orphan sweep: periodically run acquisition with no cap. In a
572+
// healthy group the sweep is a no-op — TryAcquireLease cannot
573+
// steal a valid lease — but a partition left unleased for any
574+
// reason the cap arithmetic missed (divergent heartbeat views, a
575+
// subscriber that heartbeats without acquiring) is picked up by
576+
// whichever subscriber sweeps first. An over-cap grab is shed at
577+
// the next rebalance once a peer has spare capacity to take it.
578+
uncapped := time.Since(lastOrphanSweep) >= orphanSweepInterval
579+
if uncapped {
580+
lastOrphanSweep = time.Now()
581+
}
582+
if err := s.discoverAndReconcileWorkers(ctx, sub, uncapped); err != nil {
549583
s.logger.Errorw("partition discovery failed, will retry on next tick", append(logFields, "error", err)...)
550584
}
551585
s.emitSignal(SignalPartitionUpdate)
@@ -554,8 +588,9 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
554588
}
555589

556590
// discoverAndReconcileWorkers discovers new partitions and reconciles workers.
557-
// Uses load-based fair share to limit how many partitions this subscriber acquires.
558-
func (s *subscriber) discoverAndReconcileWorkers(ctx context.Context, sub *subscription) error {
591+
// Uses fair share to limit how many partitions this subscriber acquires;
592+
// uncapped skips the fair-share cap entirely (the orphan sweep).
593+
func (s *subscriber) discoverAndReconcileWorkers(ctx context.Context, sub *subscription, uncapped bool) error {
559594
cfg := sub.config
560595

561596
// Get current leased partitions for fair share computation.
@@ -565,15 +600,21 @@ func (s *subscriber) discoverAndReconcileWorkers(ctx context.Context, sub *subsc
565600
}
566601

567602
// Use cached discovered partitions from last tick for fair share cap.
568-
// On the first tick, lastDiscoveredPartitions is nil → fairShareCap uses
569-
// only owned partitions, which gives unlimited cap for new subscribers.
603+
// On the first tick, lastDiscoveredPartitions is nil → fairShareCap sees
604+
// only owned partitions, so a joiner's first-tick cap floors at 1 and
605+
// ramps once discovery is cached.
570606
sub.workersMu.Lock()
571607
cachedDiscovered := sub.lastDiscoveredPartitions
572608
sub.workersMu.Unlock()
573609

574-
maxPartitions, err := s.fairShareCap(ctx, sub, leasedPartitions, cachedDiscovered)
575-
if err != nil {
576-
return fmt.Errorf("compute fair share cap: %w", err)
610+
// maxPartitions == 0 means unlimited (the orphan sweep, or an
611+
// uncontended single subscriber via fairShareCap).
612+
maxPartitions := 0
613+
if !uncapped {
614+
maxPartitions, err = s.fairShareCap(ctx, sub, leasedPartitions, cachedDiscovered)
615+
if err != nil {
616+
return fmt.Errorf("compute fair share cap: %w", err)
617+
}
577618
}
578619

579620
// Discover and try to acquire leases for new partitions.
@@ -1026,8 +1067,12 @@ func (s *subscriber) deregisterHeartbeat(ctx context.Context, sub *subscription)
10261067
}
10271068

10281069
// rebalance checks if this subscriber holds more partitions than its fair share
1029-
// and releases extras so other subscribers can pick them up.
1030-
func (s *subscriber) rebalance(ctx context.Context, sub *subscription, owned []string) error {
1070+
// and releases extras so other subscribers can pick them up. Returns the
1071+
// partitions actually released so the caller renews only the remainder —
1072+
// renewing a just-released lease would spuriously fail with ErrLeaseExpired.
1073+
// The owned slice is never mutated (the caller shares it with lease renewal).
1074+
// On error, partitions released before the failure are still returned.
1075+
func (s *subscriber) rebalance(ctx context.Context, sub *subscription, owned []string) (released []string, retErr error) {
10311076
cfg := sub.config
10321077

10331078
// Use cached discovered partitions from the most recent discovery tick.
@@ -1037,20 +1082,24 @@ func (s *subscriber) rebalance(ctx context.Context, sub *subscription, owned []s
10371082

10381083
maxPart, err := s.fairShareCap(ctx, sub, owned, discoveredPartitions)
10391084
if err != nil {
1040-
return fmt.Errorf("compute fair share cap: %w", err)
1085+
return nil, fmt.Errorf("compute fair share cap: %w", err)
10411086
}
10421087
if maxPart == 0 || len(owned) <= maxPart {
1043-
return nil
1088+
return nil, nil
10441089
}
10451090

1046-
// Sort deterministically so the same partitions are released across runs.
1047-
sort.Strings(owned)
1091+
// Sort a copy deterministically so the same partitions are released
1092+
// across runs without reordering the caller's slice.
1093+
sortedOwned := make([]string, len(owned))
1094+
copy(sortedOwned, owned)
1095+
sort.Strings(sortedOwned)
10481096

10491097
// Release excess partitions
1050-
for _, pk := range owned[maxPart:] {
1098+
for _, pk := range sortedOwned[maxPart:] {
10511099
if err := s.leaseStore.ReleaseLease(ctx, sub.topic, pk, cfg.SubscriberName, cfg.ConsumerGroup); err != nil {
1052-
return fmt.Errorf("release partition %s during rebalance: %w", pk, err)
1100+
return released, fmt.Errorf("release partition %s during rebalance: %w", pk, err)
10531101
}
1102+
released = append(released, pk)
10541103

10551104
// Stop the worker immediately to prevent duplicate processing.
10561105
s.stopPartitionWorker(sub, pk)
@@ -1063,14 +1112,24 @@ func (s *subscriber) rebalance(ctx context.Context, sub *subscription, owned []s
10631112
"max_partitions", maxPart,
10641113
)
10651114
}
1066-
return nil
1115+
return released, nil
10671116
}
10681117

10691118
// fairShareCap computes the max partitions this subscriber should own.
10701119
// Returns (maxPart, error). maxPart=0 means unlimited.
10711120
// owned is the caller-provided list of leased partitions.
10721121
// discoveredPartitions is an optional pre-fetched list of all known partitions;
10731122
// if nil, only owned partitions are used for fair share computation.
1123+
//
1124+
// The cap is remainder-aware: subscribers rank themselves in the sorted
1125+
// active list, the first (P mod N) ranks get floor(P/N)+1, and the rest get
1126+
// floor(P/N), so per-rank caps sum to exactly P. Independent ceil(P/N) caps
1127+
// sum to more than P and admit stable starvation states — e.g. P=12, N=5
1128+
// could settle at 3/3/3/3/0 with every subscriber at cap and nobody obliged
1129+
// to shed for the empty one. With caps summing to P, a subscriber over its
1130+
// cap implies another under its cap (rebalance sheds, the peer acquires),
1131+
// and an unleased partition implies a subscriber with spare cap to claim it
1132+
// — neither a starved subscriber nor a leftover partition is a stable state.
10741133
func (s *subscriber) fairShareCap(ctx context.Context, sub *subscription, owned []string, discoveredPartitions []string) (int, error) {
10751134
cfg := sub.config
10761135

@@ -1082,8 +1141,6 @@ func (s *subscriber) fairShareCap(ctx context.Context, sub *subscription, owned
10821141
return 0, nil
10831142
}
10841143

1085-
activeSubscribers := len(active)
1086-
10871144
// Count all known partitions as the union of owned + discovered.
10881145
// Using max(owned, discovered) would undercount when some partitions
10891146
// have leases but no messages, or vice versa.
@@ -1098,8 +1155,31 @@ func (s *subscriber) fairShareCap(ctx context.Context, sub *subscription, owned
10981155
}
10991156
totalPartitions := len(partitionSet)
11001157

1101-
// ceil(totalPartitions / activeSubscribers)
1102-
maxPart := (totalPartitions + activeSubscribers - 1) / activeSubscribers
1158+
// Rank in the sorted active list. ActiveSubscribers row order is not
1159+
// guaranteed, so sorting is what lets every subscriber derive the same
1160+
// ranking from the same set without coordination.
1161+
sort.Strings(active)
1162+
n := len(active)
1163+
rank := -1
1164+
for i, name := range active {
1165+
if name == cfg.SubscriberName {
1166+
rank = i
1167+
break
1168+
}
1169+
}
1170+
1171+
var maxPart int
1172+
if rank < 0 {
1173+
// Own heartbeat not visible this tick (e.g. the write failed): fall
1174+
// back to a conservative ceil over n+1 contenders instead of
1175+
// claiming a rank that may belong to another subscriber.
1176+
maxPart = (totalPartitions + n) / (n + 1)
1177+
} else {
1178+
maxPart = totalPartitions / n
1179+
if rank < totalPartitions%n {
1180+
maxPart++
1181+
}
1182+
}
11031183
if maxPart < 1 {
11041184
maxPart = 1
11051185
}

0 commit comments

Comments
 (0)