Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions platform/extension/messagequeue/mysql/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ subConfig.DLQ.TopicSuffix = "_dlq" // DLQ topic suffix
| `SubscriberName` | Unique worker identifier for partition leasing (e.g., hostname, pod name) |
| `ConsumerGroup` | Consumer group for independent offset tracking |
| `PollIntervalMs` | How often to poll for new messages |
| `PartitionDiscoveryIntervalMs` | How often to discover partitions, attempt lease acquisition, and reconcile workers |
| `BatchSize` | Maximum messages to fetch per poll. Set to `1` for strict serialization |
| `VisibilityTimeoutMs` | How long messages are invisible after fetch. Must exceed max processing time for `BatchSize=1` |
| `LeaseRenewalIntervalMs` | How often to renew partition leases |
Expand Down
2 changes: 1 addition & 1 deletion platform/extension/messagequeue/mysql/subscriber.go
Original file line number Diff line number Diff line change
Expand Up @@ -479,7 +479,7 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
"subscriber_name", cfg.SubscriberName,
}

discoveryTicker := time.NewTicker(time.Duration(cfg.PollIntervalMs) * time.Millisecond)
discoveryTicker := time.NewTicker(time.Duration(cfg.PartitionDiscoveryIntervalMs) * time.Millisecond)
defer discoveryTicker.Stop()

leaseTicker := time.NewTicker(time.Duration(cfg.LeaseRenewalIntervalMs) * time.Millisecond)
Expand Down
23 changes: 16 additions & 7 deletions platform/extension/messagequeue/subscription_config.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,14 @@ type SubscriptionConfig struct {
// PollIntervalMs is how often to poll for new messages (in milliseconds).
PollIntervalMs int64

// PartitionDiscoveryIntervalMs is how often to discover partitions,
// attempt lease acquisition, and reconcile partition workers (in
// milliseconds). Separate from PollIntervalMs: message polling needs low
// latency, while discovery drives topic-wide queries whose volume
// multiplies with subscribers and topics and whose outcome only changes
// on membership or partition changes.
PartitionDiscoveryIntervalMs int64

// BatchSize is the maximum number of messages to fetch per poll.
BatchSize int

Expand Down Expand Up @@ -98,13 +106,14 @@ func DLQSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionCon
// DefaultSubscriptionConfig returns a SubscriptionConfig with sensible defaults.
func DefaultSubscriptionConfig(subscriberName, consumerGroup string) SubscriptionConfig {
return SubscriptionConfig{
SubscriberName: subscriberName,
ConsumerGroup: consumerGroup,
PollIntervalMs: 100, // 100ms
BatchSize: 10,
VisibilityTimeoutMs: 60000, // 60s
LeaseRenewalIntervalMs: 10000, // 10s
LeaseDurationMs: 30000, // 30s
SubscriberName: subscriberName,
ConsumerGroup: consumerGroup,
PollIntervalMs: 100, // 100ms
PartitionDiscoveryIntervalMs: 1000, // 1s
BatchSize: 10,
VisibilityTimeoutMs: 60000, // 60s
LeaseRenewalIntervalMs: 10000, // 10s
LeaseDurationMs: 30000, // 30s
Retry: RetryConfig{
MaxAttempts: 3,
InitialBackoffMs: 1000, // 1s
Expand Down
30 changes: 30 additions & 0 deletions test/integration/extension/messagequeue/mysql/queue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,8 +102,12 @@ func (s *SQLQueueIntegrationSuite) TearDownSuite() {
// timeouts for fast integration tests. The defaults (30s lease, 60s visibility)
// would make crash recovery tests wait 90s of real wall-clock time since the
// subscriber can't find invisible messages until the DB timeout expires.
// Partition discovery is likewise pinned to 100ms so initial lease
// acquisition and rebalance convergence stay fast under the 1s production
// default.
func testSubConfig(subscriberName, consumerGroup string) extqueue.SubscriptionConfig {
cfg := extqueue.DefaultSubscriptionConfig(subscriberName, consumerGroup)
cfg.PartitionDiscoveryIntervalMs = 100
cfg.VisibilityTimeoutMs = 2000
cfg.LeaseDurationMs = 3000
cfg.LeaseRenewalIntervalMs = 1000
Expand Down Expand Up @@ -344,6 +348,7 @@ func (s *SQLQueueIntegrationSuite) TestPublishAndSubscribe() {

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

subConfig := extqueue.DefaultSubscriptionConfig("reset-worker", consumerGroup)
subConfig.PartitionDiscoveryIntervalMs = 100
subConfig.PollIntervalMs = 100
deliveryChan, err := subscriber.Subscribe(s.ctx, topic, subConfig)
require.NoError(t, err)
Expand Down Expand Up @@ -2115,6 +2137,7 @@ func (s *SQLQueueIntegrationSuite) TestInFlightMessageDoesNotBlockOtherMessages(
// Subscribe with batch=10 to fetch multiple messages per poll. The default
// 60s visibility timeout keeps msg-1 invisible for the whole test.
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "nack-nb-cg")
subConfig.PartitionDiscoveryIntervalMs = 100
subConfig.PollIntervalMs = 50
subConfig.BatchSize = 10
deliveryCh, err := q.Subscriber().Subscribe(s.ctx, topic, subConfig)
Expand Down Expand Up @@ -2172,6 +2195,7 @@ func (s *SQLQueueIntegrationSuite) TestPostponeBlocksPartitionUntilDue() {
// Subscribe with batch=10 so the barrier — not the batch size — is what
// keeps later offsets back.
subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "postpone-cg")
subConfig.PartitionDiscoveryIntervalMs = 100
subConfig.PollIntervalMs = 50
subConfig.BatchSize = 10
deliveryCh, err := q.Subscriber().Subscribe(s.ctx, topic, subConfig)
Expand Down Expand Up @@ -2268,6 +2292,7 @@ func (s *SQLQueueIntegrationSuite) TestPostponeResetsRetryBudget() {

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

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

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

// Two consumer groups subscribing to the same topic
cfg1 := extqueue.DefaultSubscriptionConfig("worker-1", "cg-alpha")
cfg1.PartitionDiscoveryIntervalMs = 100
cfg1.PollIntervalMs = 50
cfg2 := extqueue.DefaultSubscriptionConfig("worker-2", "cg-beta")
cfg2.PartitionDiscoveryIntervalMs = 100
cfg2.PollIntervalMs = 50

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

subConfig := extqueue.DefaultSubscriptionConfig("worker-1", "watermark-cg")
subConfig.PartitionDiscoveryIntervalMs = 100
subConfig.PollIntervalMs = 100
subConfig.VisibilityTimeoutMs = 30000 // long visibility so nothing re-delivers
subConfig.BatchSize = 10
Expand Down
Loading