Skip to content

Commit c26d89e

Browse files
committed
fix(messagequeue): stop unbounded heartbeat row growth
## Summary ### Why? `queue_subscriber_heartbeats` rows are never removed: `Deregister` only soft-deletes (sets `deregistered_at`), crashed subscribers never deregister at all, and every process registers under a fresh `hostname-pid` subscriber name — so the table grows monotonically with nodes × subscriptions × deploys. `ActiveSubscribers` range-scans all rows for `(consumer_group, topic)`, so the dead rows slowly tax every fair-share computation. ### What? - `Deregister` now hard-deletes the row. Subscriber names are unique per process, so a departed subscriber's row has no further use; re-subscribing re-inserts via the heartbeat upsert. - New `PurgeStale` store method deletes rows whose heartbeat is older than a threshold; the subscriber calls it each lease tick with 10x `LeaseDurationMs` (5min at defaults) as the backstop for subscribers that crashed without deregistering. Purging a live-but-stalled subscriber's row is harmless — its next heartbeat re-inserts it. - The `deregistered_at` column and its query filter are retained (no schema change) for compatibility with rows written by older code. ## Test Plan - ✅ Integration: `TestRebalance_SubscriberLeaves` now also asserts the closed subscriber's heartbeat row is gone from the table.
1 parent 474e9c9 commit c26d89e

7 files changed

Lines changed: 152 additions & 23 deletions

File tree

platform/extension/messagequeue/mysql/mock_stores.go

Lines changed: 14 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

platform/extension/messagequeue/mysql/stores.go

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -154,8 +154,18 @@ type subscriberHeartbeatStore interface {
154154
// within this duration are considered dead.
155155
ActiveSubscribers(ctx context.Context, topic string, consumerGroup string, staleDurationMs int64) ([]string, error)
156156

157-
// Deregister removes a subscriber's heartbeat entry
157+
// Deregister removes a subscriber's heartbeat row. Hard delete: the row
158+
// is not needed once the subscriber is gone, and subscriber names are
159+
// unique per process (hostname-pid), so rows would otherwise accumulate
160+
// forever across deploys. Re-subscribing re-inserts via Heartbeat.
158161
Deregister(ctx context.Context, topic string, subscriberName string, consumerGroup string) error
162+
163+
// PurgeStale deletes heartbeat rows whose last heartbeat is older than
164+
// olderThanMs. Backstop for subscribers that never deregistered
165+
// (crashes, SIGKILL): without it the table grows monotonically since
166+
// every process registers under a fresh name. Deleting a live-but-stalled
167+
// subscriber's row is harmless — its next heartbeat re-inserts it.
168+
PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) error
159169
}
160170

161171
// DeliveryState represents the full per-message delivery tracking state.

platform/extension/messagequeue/mysql/subscriber.go

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,14 @@ const (
5353
// queries when many partitions are idle (e.g., 50 idle partitions at 100ms
5454
// poll interval = 500 GC queries/sec without throttling).
5555
gcIdleTickInterval = 100
56+
57+
// heartbeatPurgeAfterLeaseDurations sets the age threshold for purging
58+
// abandoned heartbeat rows, as a multiple of LeaseDurationMs (10x = 5min
59+
// at defaults). Well past every transient window in the protocol — a row
60+
// that stale belongs to a subscriber that crashed without deregistering.
61+
// Purging a live-but-stalled subscriber's row is harmless: its next
62+
// heartbeat re-inserts it.
63+
heartbeatPurgeAfterLeaseDurations = 10
5664
)
5765

5866
// HookSignal identifies the type of subscriber lifecycle event.
@@ -565,6 +573,13 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
565573
if err := s.sendHeartbeat(ctx, sub); err != nil {
566574
s.logger.Errorw("periodic heartbeat failed", append(logFields, "error", err)...)
567575
}
576+
// Purge heartbeat rows abandoned by subscribers that never
577+
// deregistered (crashes) — without this the table grows
578+
// monotonically, since every process registers under a fresh
579+
// hostname-pid name.
580+
if err := s.heartbeatStore.PurgeStale(ctx, sub.topic, cfg.ConsumerGroup, heartbeatPurgeAfterLeaseDurations*cfg.LeaseDurationMs); err != nil {
581+
s.logger.Errorw("stale heartbeat purge failed", append(logFields, "error", err)...)
582+
}
568583
s.emitSignal(SignalPartitionUpdate)
569584

570585
case <-discoveryTicker.C:

platform/extension/messagequeue/mysql/subscriber_heartbeat_store.go

Lines changed: 35 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -101,18 +101,16 @@ func (s *sqlSubscriberHeartbeatStore) ActiveSubscribers(ctx context.Context, top
101101
return names, nil
102102
}
103103

104-
// Deregister soft-deletes a subscriber's heartbeat entry by setting deregistered_at.
105-
// Idempotent: no-op if already deregistered.
104+
// Deregister removes a subscriber's heartbeat row (hard delete — see the
105+
// subscriberHeartbeatStore interface doc). Idempotent: no-op if already gone.
106106
func (s *sqlSubscriberHeartbeatStore) Deregister(ctx context.Context, topic string, subscriberName string, consumerGroup string) (retErr error) {
107107
op := metrics.Begin(s.scope, "deregister", metrics.StorageLatencyBuckets)
108108
defer func() { op.Complete(retErr) }()
109109

110-
now := s.nowFunc().UnixMilli()
111-
112110
_, err := s.db.ExecContext(ctx, fmt.Sprintf(`
113-
UPDATE %s SET deregistered_at = ?
114-
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ? AND deregistered_at = 0
115-
`, SubscriberHeartbeatsTableName), now, consumerGroup, topic, subscriberName)
111+
DELETE FROM %s
112+
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ?
113+
`, SubscriberHeartbeatsTableName), consumerGroup, topic, subscriberName)
116114

117115
if err != nil {
118116
return fmt.Errorf("failed to deregister subscriber: %w", err)
@@ -125,3 +123,33 @@ func (s *sqlSubscriberHeartbeatStore) Deregister(ctx context.Context, topic stri
125123

126124
return nil
127125
}
126+
127+
// PurgeStale deletes heartbeat rows older than olderThanMs for the topic and
128+
// consumer group. See the subscriberHeartbeatStore interface doc.
129+
func (s *sqlSubscriberHeartbeatStore) PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) (retErr error) {
130+
op := metrics.Begin(s.scope, "purge_stale", metrics.StorageLatencyBuckets)
131+
defer func() { op.Complete(retErr) }()
132+
133+
threshold := s.nowFunc().UnixMilli() - olderThanMs
134+
135+
result, err := s.db.ExecContext(ctx, fmt.Sprintf(`
136+
DELETE FROM %s
137+
WHERE consumer_group = ? AND topic = ? AND heartbeat_at < ?
138+
`, SubscriberHeartbeatsTableName), consumerGroup, topic, threshold)
139+
140+
if err != nil {
141+
return fmt.Errorf("failed to purge stale heartbeats: %w", err)
142+
}
143+
144+
// RowsAffected error is swallowed because the DELETE itself succeeded;
145+
// the count is for observability only.
146+
if deleted, err := result.RowsAffected(); err == nil && deleted > 0 {
147+
metrics.NamedCounter(s.scope, "purge_stale", "rows_deleted", deleted, metrics.NewTag("topic", topic))
148+
s.logger.Debugw("purged stale heartbeats",
149+
logTopic, topic,
150+
"deleted", deleted,
151+
)
152+
}
153+
154+
return nil
155+
}

platform/extension/messagequeue/mysql/subscriber_heartbeat_store_test.go

Lines changed: 67 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -172,22 +172,74 @@ func TestSubscriberHeartbeatStore_ActiveSubscribers_ExcludesDeregistered(t *test
172172
require.NoError(t, mock.ExpectationsWereMet())
173173
}
174174

175-
func TestSubscriberHeartbeatStore_Deregister_SoftDelete(t *testing.T) {
175+
func TestSubscriberHeartbeatStore_Deregister_HardDelete(t *testing.T) {
176176
db, mock, store := setupSubscriberHeartbeatStoreTest(t)
177177
defer db.Close()
178178

179179
ctx := context.Background()
180180

181-
// Verify deregister uses UPDATE (not DELETE) and targets only active rows (deregistered_at = 0)
182-
mock.ExpectExec(`UPDATE queue_subscriber_heartbeats SET deregistered_at.*AND deregistered_at = 0`).
183-
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
181+
// Verify deregister deletes the row outright — subscriber names are
182+
// unique per process, so soft-deleted rows would accumulate forever.
183+
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
184+
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
184185
WillReturnResult(sqlmock.NewResult(0, 1))
185186

186187
err := store.Deregister(ctx, "test_topic", testSubscriberName, testConsumerGroup)
187188
require.NoError(t, err)
188189
require.NoError(t, mock.ExpectationsWereMet())
189190
}
190191

192+
func TestSubscriberHeartbeatStore_PurgeStale(t *testing.T) {
193+
tests := []struct {
194+
name string
195+
setup func(mock sqlmock.Sqlmock)
196+
wantErr bool
197+
}{
198+
{
199+
name: "deletes rows older than threshold",
200+
setup: func(mock sqlmock.Sqlmock) {
201+
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
202+
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
203+
WillReturnResult(sqlmock.NewResult(0, 3))
204+
},
205+
},
206+
{
207+
name: "no stale rows is a no-op",
208+
setup: func(mock sqlmock.Sqlmock) {
209+
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
210+
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
211+
WillReturnResult(sqlmock.NewResult(0, 0))
212+
},
213+
},
214+
{
215+
name: "database error",
216+
setup: func(mock sqlmock.Sqlmock) {
217+
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
218+
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
219+
WillReturnError(fmt.Errorf("db error"))
220+
},
221+
wantErr: true,
222+
},
223+
}
224+
225+
for _, tt := range tests {
226+
t.Run(tt.name, func(t *testing.T) {
227+
db, mock, store := setupSubscriberHeartbeatStoreTest(t)
228+
defer db.Close()
229+
230+
tt.setup(mock)
231+
232+
err := store.PurgeStale(context.Background(), "test_topic", testConsumerGroup, 300_000)
233+
if tt.wantErr {
234+
require.Error(t, err)
235+
} else {
236+
require.NoError(t, err)
237+
}
238+
require.NoError(t, mock.ExpectationsWereMet())
239+
})
240+
}
241+
}
242+
191243
func TestSubscriberHeartbeatStore_ReRegistration(t *testing.T) {
192244
db, mock, store := setupSubscriberHeartbeatStoreTest(t)
193245
defer db.Close()
@@ -199,15 +251,15 @@ func TestSubscriberHeartbeatStore_ReRegistration(t *testing.T) {
199251
WithArgs(testConsumerGroup, "test_topic", testSubscriberName, sqlmock.AnyArg()).
200252
WillReturnResult(sqlmock.NewResult(1, 1))
201253

202-
// Step 2: Deregister soft-deletes the subscriber
203-
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
204-
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
254+
// Step 2: Deregister deletes the subscriber's row
255+
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
256+
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
205257
WillReturnResult(sqlmock.NewResult(0, 1))
206258

207-
// Step 3: Heartbeat again re-registers (ON DUPLICATE KEY UPDATE resets deregistered_at = 0)
259+
// Step 3: Heartbeat again re-registers with a fresh insert
208260
mock.ExpectExec("INSERT INTO queue_subscriber_heartbeats").
209261
WithArgs(testConsumerGroup, "test_topic", testSubscriberName, sqlmock.AnyArg()).
210-
WillReturnResult(sqlmock.NewResult(0, 2)) // 2 = ON DUPLICATE KEY UPDATE
262+
WillReturnResult(sqlmock.NewResult(1, 1))
211263

212264
err := store.Heartbeat(ctx, "test_topic", testSubscriberName, testConsumerGroup)
213265
require.NoError(t, err)
@@ -230,26 +282,26 @@ func TestSubscriberHeartbeatStore_Deregister(t *testing.T) {
230282
{
231283
name: "successfully deregister",
232284
setup: func(mock sqlmock.Sqlmock) {
233-
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
234-
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
285+
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
286+
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
235287
WillReturnResult(sqlmock.NewResult(0, 1))
236288
},
237289
wantErr: false,
238290
},
239291
{
240292
name: "idempotent - already deregistered",
241293
setup: func(mock sqlmock.Sqlmock) {
242-
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
243-
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
294+
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
295+
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
244296
WillReturnResult(sqlmock.NewResult(0, 0))
245297
},
246298
wantErr: false,
247299
},
248300
{
249301
name: "database error",
250302
setup: func(mock sqlmock.Sqlmock) {
251-
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
252-
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
303+
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
304+
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
253305
WillReturnError(fmt.Errorf("db error"))
254306
},
255307
wantErr: true,

platform/extension/messagequeue/mysql/subscriber_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ func newTestHeartbeatStore(ctrl *gomock.Controller) *MocksubscriberHeartbeatStor
4141
mockHB.EXPECT().Heartbeat(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
4242
mockHB.EXPECT().ActiveSubscribers(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return([]string{"self"}, nil).AnyTimes()
4343
mockHB.EXPECT().Deregister(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
44+
mockHB.EXPECT().PurgeStale(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
4445
return mockHB
4546
}
4647

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1814,6 +1814,15 @@ func (s *SQLQueueIntegrationSuite) TestRebalance_SubscriberLeaves() {
18141814
return len(leases["s1"]) == 4
18151815
}, "S1 should reacquire all 4 partitions after S2 leaves")
18161816

1817+
// Deregistration hard-deletes the heartbeat row — subscriber names are
1818+
// unique per process, so rows would otherwise accumulate forever.
1819+
var s2Rows int
1820+
require.NoError(t, s.db.QueryRowContext(s.ctx, `
1821+
SELECT COUNT(*) FROM queue_subscriber_heartbeats
1822+
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ?
1823+
`, consumerGroup, topic, "s2").Scan(&s2Rows))
1824+
assert.Equal(t, 0, s2Rows, "closed subscriber's heartbeat row must be deleted")
1825+
18171826
t.Logf("Subscriber leave verified: S1 owns all 4 partitions after S2 departed")
18181827
}
18191828

0 commit comments

Comments
 (0)