Skip to content

Commit 4f616b7

Browse files
committed
feat(messagequeue): split partition discovery cadence from polling
## Summary ### Why? Partition discovery, lease acquisition, and worker reconciliation all run on the message-poll ticker (`PollIntervalMs`, 100ms default). Discovery drives topic-wide work every tick — a `DISTINCT partition_key` scan, an active-subscriber read, self-lease reads, and lease-acquisition probes — whose query volume multiplies with subscribers × topics at 10x/sec, even though its outcome only changes when membership or the partition set changes. Message polling needs 100ms latency; discovery does not. ### What? New `PartitionDiscoveryIntervalMs` subscription config (default 1s) drives the supervisor's discovery ticker; `PollIntervalMs` continues to drive per-partition message polling unchanged. This cuts discovery-driven query volume ~10x at the default settings. The accepted trade-off is that a brand-new partition's first message now waits up to the discovery interval (1s) before a worker picks it up; messages on already-owned partitions are unaffected. Integration test configs pin discovery to 100ms so lease-handoff and rebalance convergence assertions stay fast. ## Test Plan - ✅ Full Docker integration suite (`bazel test //test/integration/extension/messagequeue/...`) — discovery-latency-sensitive tests (empty-topic wake-up, rebalance convergence, crash recovery) pass with the new cadence.
1 parent 1db0efc commit 4f616b7

4 files changed

Lines changed: 48 additions & 8 deletions

File tree

platform/extension/messagequeue/mysql/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ subConfig.DLQ.TopicSuffix = "_dlq" // DLQ topic suffix
6969
| `SubscriberName` | Unique worker identifier for partition leasing (e.g., hostname, pod name) |
7070
| `ConsumerGroup` | Consumer group for independent offset tracking |
7171
| `PollIntervalMs` | How often to poll for new messages |
72+
| `PartitionDiscoveryIntervalMs` | How often to discover partitions, attempt lease acquisition, and reconcile workers |
7273
| `BatchSize` | Maximum messages to fetch per poll. Set to `1` for strict serialization |
7374
| `VisibilityTimeoutMs` | How long messages are invisible after fetch. Must exceed max processing time for `BatchSize=1` |
7475
| `LeaseRenewalIntervalMs` | How often to renew partition leases |

platform/extension/messagequeue/mysql/subscriber.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -479,7 +479,7 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
479479
"subscriber_name", cfg.SubscriberName,
480480
}
481481

482-
discoveryTicker := time.NewTicker(time.Duration(cfg.PollIntervalMs) * time.Millisecond)
482+
discoveryTicker := time.NewTicker(time.Duration(cfg.PartitionDiscoveryIntervalMs) * time.Millisecond)
483483
defer discoveryTicker.Stop()
484484

485485
leaseTicker := time.NewTicker(time.Duration(cfg.LeaseRenewalIntervalMs) * time.Millisecond)

platform/extension/messagequeue/subscription_config.go

Lines changed: 16 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,14 @@ type SubscriptionConfig struct {
3030
// PollIntervalMs is how often to poll for new messages (in milliseconds).
3131
PollIntervalMs int64
3232

33+
// PartitionDiscoveryIntervalMs is how often to discover partitions,
34+
// attempt lease acquisition, and reconcile partition workers (in
35+
// milliseconds). Separate from PollIntervalMs: message polling needs low
36+
// latency, while discovery drives topic-wide queries whose volume
37+
// multiplies with subscribers and topics and whose outcome only changes
38+
// on membership or partition changes.
39+
PartitionDiscoveryIntervalMs int64
40+
3341
// BatchSize is the maximum number of messages to fetch per poll.
3442
BatchSize int
3543

@@ -98,13 +106,14 @@ func DLQSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionCon
98106
// DefaultSubscriptionConfig returns a SubscriptionConfig with sensible defaults.
99107
func DefaultSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionConfig {
100108
return SubscriptionConfig{
101-
SubscriberName: subscriberName,
102-
ConsumerGroup: consumerGroup,
103-
PollIntervalMs: 100, // 100ms
104-
BatchSize: 10,
105-
VisibilityTimeoutMs: 60000, // 60s
106-
LeaseRenewalIntervalMs: 10000, // 10s
107-
LeaseDurationMs: 30000, // 30s
109+
SubscriberName: subscriberName,
110+
ConsumerGroup: consumerGroup,
111+
PollIntervalMs: 100, // 100ms
112+
PartitionDiscoveryIntervalMs: 1000, // 1s
113+
BatchSize: 10,
114+
VisibilityTimeoutMs: 60000, // 60s
115+
LeaseRenewalIntervalMs: 10000, // 10s
116+
LeaseDurationMs: 30000, // 30s
108117
Retry: RetryConfig{
109118
MaxAttempts: 3,
110119
InitialBackoffMs: 1000, // 1s

test/integration/extension/messagequeue/mysql/queue_test.go

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,8 +102,12 @@ func (s *SQLQueueIntegrationSuite) TearDownSuite() {
102102
// timeouts for fast integration tests. The defaults (30s lease, 60s visibility)
103103
// would make crash recovery tests wait 90s of real wall-clock time since the
104104
// subscriber can't find invisible messages until the DB timeout expires.
105+
// Partition discovery is likewise pinned to 100ms so initial lease
106+
// acquisition and rebalance convergence stay fast under the 1s production
107+
// default.
105108
func testSubConfig(subscriberName, consumerGroup string) extqueue.SubscriptionConfig {
106109
cfg := extqueue.DefaultSubscriptionConfig(subscriberName, consumerGroup)
110+
cfg.PartitionDiscoveryIntervalMs = 100
107111
cfg.VisibilityTimeoutMs = 2000
108112
cfg.LeaseDurationMs = 3000
109113
cfg.LeaseRenewalIntervalMs = 1000
@@ -344,6 +348,7 @@ func (s *SQLQueueIntegrationSuite) TestPublishAndSubscribe() {
344348

345349
// Subscribe first with config
346350
subConfig := extqueue.DefaultSubscriptionConfig("test-worker-1", "test-consumer")
351+
subConfig.PartitionDiscoveryIntervalMs = 100
347352
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
348353
require.NoError(t, err)
349354

@@ -414,6 +419,7 @@ func (s *SQLQueueIntegrationSuite) TestSubscriberPerPartitionIsolation() {
414419

415420
// Subscribe with short poll interval for fast test
416421
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "isolation-consumer")
422+
subConfig.PartitionDiscoveryIntervalMs = 100
417423
subConfig.PollIntervalMs = 100
418424
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
419425
require.NoError(t, err)
@@ -481,6 +487,7 @@ func (s *SQLQueueIntegrationSuite) TestSubscriberPartitionOrderPreserved() {
481487

482488
// Subscribe and receive all
483489
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "order-consumer")
490+
subConfig.PartitionDiscoveryIntervalMs = 100
484491
subConfig.PollIntervalMs = 100
485492
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
486493
require.NoError(t, err)
@@ -521,6 +528,7 @@ func (s *SQLQueueIntegrationSuite) TestMultiplePartitions() {
521528

522529
// Subscribe
523530
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "multi-partition-consumer")
531+
subConfig.PartitionDiscoveryIntervalMs = 100
524532
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
525533
require.NoError(t, err)
526534

@@ -650,6 +658,7 @@ func (s *SQLQueueIntegrationSuite) TestIdempotentPublish() {
650658

651659
// Subscribe
652660
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "idempotent-consumer")
661+
subConfig.PartitionDiscoveryIntervalMs = 100
653662
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
654663
require.NoError(t, err)
655664

@@ -696,6 +705,7 @@ func (s *SQLQueueIntegrationSuite) TestConcurrentPublishers() {
696705

697706
// Subscribe
698707
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "concurrent-consumer")
708+
subConfig.PartitionDiscoveryIntervalMs = 100
699709
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
700710
require.NoError(t, err)
701711

@@ -825,10 +835,12 @@ func (s *SQLQueueIntegrationSuite) TestMultipleConsumerGroups() {
825835

826836
// Subscribe both groups
827837
subConfig1 := extqueue.DefaultSubscriptionConfig("worker-1", "group-A")
838+
subConfig1.PartitionDiscoveryIntervalMs = 100
828839
deliveryChan1, err := subscriber1.Subscribe(s.ctx, topic, subConfig1)
829840
require.NoError(t, err)
830841

831842
subConfig2 := extqueue.DefaultSubscriptionConfig("worker-1", "group-B")
843+
subConfig2.PartitionDiscoveryIntervalMs = 100
832844
deliveryChan2, err := subscriber2.Subscribe(s.ctx, topic, subConfig2)
833845
require.NoError(t, err)
834846

@@ -904,10 +916,12 @@ func (s *SQLQueueIntegrationSuite) TestMultipleWorkersInConsumerGroup() {
904916

905917
// Subscribe both workers
906918
subConfig1 := extqueue.DefaultSubscriptionConfig("worker-1", consumerGroup)
919+
subConfig1.PartitionDiscoveryIntervalMs = 100
907920
deliveryChan1, err := subscriber1.Subscribe(s.ctx, topic, subConfig1)
908921
require.NoError(t, err)
909922

910923
subConfig2 := extqueue.DefaultSubscriptionConfig("worker-2", consumerGroup)
924+
subConfig2.PartitionDiscoveryIntervalMs = 100
911925
deliveryChan2, err := subscriber2.Subscribe(s.ctx, topic, subConfig2)
912926
require.NoError(t, err)
913927

@@ -979,6 +993,7 @@ func (s *SQLQueueIntegrationSuite) TestConcurrentSubscribers() {
979993

980994
subscriber := q.Subscriber()
981995
subConfig := extqueue.DefaultSubscriptionConfig(fmt.Sprintf("worker-%d", i), consumerGroup)
996+
subConfig.PartitionDiscoveryIntervalMs = 100
982997
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
983998
require.NoError(t, err)
984999
deliveryChans = append(deliveryChans, deliveryChan)
@@ -1078,6 +1093,7 @@ func (s *SQLQueueIntegrationSuite) TestDeadLetterQueue() {
10781093
t.Logf("Subscribing to DLQ topic: %s", dlqTopic)
10791094

10801095
dlqConfig := extqueue.DefaultSubscriptionConfig("worker-1", "dlq-consumer")
1096+
dlqConfig.PartitionDiscoveryIntervalMs = 100
10811097
dlqDeliveryChan, err := subscriber.Subscribe(s.ctx, dlqTopic, dlqConfig)
10821098
require.NoError(t, err)
10831099

@@ -1129,6 +1145,7 @@ func (s *SQLQueueIntegrationSuite) TestMessageOrderingWithinPartition() {
11291145

11301146
// Subscribe first
11311147
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "ordering-consumer")
1148+
subConfig.PartitionDiscoveryIntervalMs = 100
11321149
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
11331150
require.NoError(t, err)
11341151

@@ -1191,6 +1208,7 @@ func (s *SQLQueueIntegrationSuite) TestLateSubscriber() {
11911208
// Now subscribe (late subscriber)
11921209
subscriber := q.Subscriber()
11931210
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "late-consumer")
1211+
subConfig.PartitionDiscoveryIntervalMs = 100
11941212
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
11951213
require.NoError(t, err)
11961214
t.Logf("Late subscriber joined after messages published")
@@ -1232,6 +1250,7 @@ func (s *SQLQueueIntegrationSuite) TestEmptyTopicSubscribe() {
12321250

12331251
// Subscribe to empty topic (no messages published yet)
12341252
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "empty-consumer")
1253+
subConfig.PartitionDiscoveryIntervalMs = 100
12351254
subConfig.PollIntervalMs = 100 // 100 milliseconds
12361255
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
12371256
require.NoError(t, err)
@@ -1490,6 +1509,7 @@ func (s *SQLQueueIntegrationSuite) TestAdmin_ConsumerLagAfterPartialAck() {
14901509

14911510
// Subscribe and ack only 2
14921511
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", consumerGroup)
1512+
subConfig.PartitionDiscoveryIntervalMs = 100
14931513
subConfig.PollIntervalMs = 100
14941514
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
14951515
require.NoError(t, err)
@@ -1539,6 +1559,7 @@ func (s *SQLQueueIntegrationSuite) TestAdmin_LeasesAndOffsets() {
15391559
require.NoError(t, publisher.Publish(s.ctx, topic, entityqueue.NewMessage("lo-1", []byte("a"), "p1", nil)))
15401560

15411561
subConfig := extqueue.DefaultSubscriptionConfig("admin-worker-1", consumerGroup)
1562+
subConfig.PartitionDiscoveryIntervalMs = 100
15421563
subConfig.PollIntervalMs = 100
15431564
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
15441565
require.NoError(t, err)
@@ -1615,6 +1636,7 @@ func (s *SQLQueueIntegrationSuite) TestAdmin_ResetOffsetAndReleaseLease() {
16151636
require.NoError(t, publisher.Publish(s.ctx, topic, entityqueue.NewMessage("r1", []byte("a"), "rp1", nil)))
16161637

16171638
subConfig := extqueue.DefaultSubscriptionConfig("reset-worker", consumerGroup)
1639+
subConfig.PartitionDiscoveryIntervalMs = 100
16181640
subConfig.PollIntervalMs = 100
16191641
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
16201642
require.NoError(t, err)
@@ -2115,6 +2137,7 @@ func (s *SQLQueueIntegrationSuite) TestInFlightMessageDoesNotBlockOtherMessages(
21152137
// Subscribe with batch=10 to fetch multiple messages per poll. The default
21162138
// 60s visibility timeout keeps msg-1 invisible for the whole test.
21172139
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "nack-nb-cg")
2140+
subConfig.PartitionDiscoveryIntervalMs = 100
21182141
subConfig.PollIntervalMs = 50
21192142
subConfig.BatchSize = 10
21202143
deliveryCh, err := q.Subscriber().Subscribe(s.ctx, topic, subConfig)
@@ -2172,6 +2195,7 @@ func (s *SQLQueueIntegrationSuite) TestPostponeBlocksPartitionUntilDue() {
21722195
// Subscribe with batch=10 so the barrier — not the batch size — is what
21732196
// keeps later offsets back.
21742197
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "postpone-cg")
2198+
subConfig.PartitionDiscoveryIntervalMs = 100
21752199
subConfig.PollIntervalMs = 50
21762200
subConfig.BatchSize = 10
21772201
deliveryCh, err := q.Subscriber().Subscribe(s.ctx, topic, subConfig)
@@ -2268,6 +2292,7 @@ func (s *SQLQueueIntegrationSuite) TestPostponeResetsRetryBudget() {
22682292

22692293
dlqTopic := topic + subConfig.DLQ.TopicSuffix
22702294
dlqConfig := extqueue.DefaultSubscriptionConfig("worker-1", "postpone-budget-cg")
2295+
dlqConfig.PartitionDiscoveryIntervalMs = 100
22712296
dlqDeliveryChan, err := q.Subscriber().Subscribe(s.ctx, dlqTopic, dlqConfig)
22722297
require.NoError(t, err)
22732298

@@ -2299,6 +2324,7 @@ func (s *SQLQueueIntegrationSuite) TestBatchSizeOneStrictSerialization() {
22992324

23002325
// Subscribe with batchSize=1 for strict serialization
23012326
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "serial-cg")
2327+
subConfig.PartitionDiscoveryIntervalMs = 100
23022328
subConfig.PollIntervalMs = 50
23032329
subConfig.BatchSize = 1
23042330
deliveryCh, err := q.Subscriber().Subscribe(s.ctx, topic, subConfig)
@@ -2346,8 +2372,10 @@ func (s *SQLQueueIntegrationSuite) TestMultipleConsumerGroupsIndependentState()
23462372

23472373
// Two consumer groups subscribing to the same topic
23482374
cfg1 := extqueue.DefaultSubscriptionConfig("worker-1", "cg-alpha")
2375+
cfg1.PartitionDiscoveryIntervalMs = 100
23492376
cfg1.PollIntervalMs = 50
23502377
cfg2 := extqueue.DefaultSubscriptionConfig("worker-2", "cg-beta")
2378+
cfg2.PartitionDiscoveryIntervalMs = 100
23512379
cfg2.PollIntervalMs = 50
23522380

23532381
ch1, err := q.Subscriber().Subscribe(s.ctx, topic, cfg1)
@@ -2480,6 +2508,7 @@ func (s *SQLQueueIntegrationSuite) TestCrashAfterRejectDoesNotLoseMessages() {
24802508
// Verify DLQ contains msg-B
24812509
dlqTopic := topic + subConfig.DLQ.TopicSuffix
24822510
dlqConfig := extqueue.DefaultSubscriptionConfig("worker-2", "crash-reject-cg")
2511+
dlqConfig.PartitionDiscoveryIntervalMs = 100
24832512
dlqConfig.PollIntervalMs = 100
24842513
dlqChan, err := q2.Subscriber().Subscribe(s.ctx, dlqTopic, dlqConfig)
24852514
require.NoError(t, err)
@@ -2637,6 +2666,7 @@ func (s *SQLQueueIntegrationSuite) TestWatermarkAdvancesContiguously() {
26372666
}
26382667

26392668
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "watermark-cg")
2669+
subConfig.PartitionDiscoveryIntervalMs = 100
26402670
subConfig.PollIntervalMs = 100
26412671
subConfig.VisibilityTimeoutMs = 30000 // long visibility so nothing re-delivers
26422672
subConfig.BatchSize = 10

0 commit comments

Comments
 (0)