Skip to content

Commit 0a35fc0

Browse files
authored
Merge branch 'main' into preetam/partition-by-queue
2 parents 253fb33 + 917b3d4 commit 0a35fc0

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)