Skip to content

Commit 503e718

Browse files
authored
[server] Fix follower heartbeat entries for AA stores to use local region only (linkedin#2584)
Fix AA follower heartbeat initialization to track only local region Problem: Active/Active follower replicas were initializing heartbeat tracking for all regions instead of only the local region. Because followers consume only the local VT, non-local heartbeat entries were never updated, causing unbounded delay growth and misleading catch-up metrics for remote regions. Solution: Restore the leader-versus-follower initialization logic in HeartbeatMonitoringService so AA leaders still track all regions, while AA followers create heartbeat entries only for the local region. Update HeartbeatMonitoringServiceTest to verify followers keep a single local-region entry and do not create remote-region entries.
1 parent 97e3800 commit 503e718

3 files changed

Lines changed: 7 additions & 9 deletions

File tree

clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/ingestion/heartbeat/HeartbeatMonitoringService.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -133,7 +133,7 @@ private synchronized void initializeEntry(
133133
String storeName = version.getStoreName();
134134
int versionNum = version.getNumber();
135135
long currentTime = System.currentTimeMillis();
136-
if (version.isActiveActiveReplicationEnabled()) {
136+
if (version.isActiveActiveReplicationEnabled() && !isFollower) {
137137
for (String region: regionNames) {
138138
if (Utils.isSeparateTopicRegion(region) && !version.isSeparateRealTimeTopicEnabled()) {
139139
continue;

clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ingestion/heartbeat/HeartbeatMonitoringServiceTest.java

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -413,10 +413,8 @@ public void testAddLeaderLagMonitor(boolean enableSepRT) {
413413

414414
// Version 3 (futureVersion) is non-A/A, so followers only have the local region
415415
Assert.assertEquals(countRegions(heartbeatMonitoringService.getFollowerHeartbeatTimeStamps(), TEST_STORE, 3, 0), 1);
416-
// Version 2 (currentVersion) is A/A, so followers have all regions initialized
417-
Assert.assertEquals(
418-
countRegions(heartbeatMonitoringService.getFollowerHeartbeatTimeStamps(), TEST_STORE, 2, 0),
419-
2 + (enableSepRT ? 1 : 0));
416+
// Version 2 (currentVersion) is A/A, but followers should only have local region
417+
Assert.assertEquals(countRegions(heartbeatMonitoringService.getFollowerHeartbeatTimeStamps(), TEST_STORE, 2, 0), 1);
420418

421419
// make sure we didn't get any leader heartbeats yet recorded
422420
Assert.assertFalse(hasStore(heartbeatMonitoringService.getLeaderHeartbeatTimeStamps(), TEST_STORE));
@@ -525,9 +523,8 @@ public void testAddLeaderLagMonitor(boolean enableSepRT) {
525523
futureVersion.getNumber(),
526524
1,
527525
REMOTE_FABRIC));
528-
// currentVersion (A/A): REMOTE_FABRIC entry should exist for followers since entries are initialized for all
529-
// regions
530-
Assert.assertNotNull(
526+
// currentVersion (A/A): REMOTE_FABRIC entry should NOT exist for followers since only local region is initialized
527+
Assert.assertNull(
531528
getEntry(
532529
heartbeatMonitoringService.getFollowerHeartbeatTimeStamps(),
533530
TEST_STORE,

internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/endToEnd/TestSeparateRealtimeTopicIngestion.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -315,7 +315,8 @@ private void validateSeparateRealtimeTopicHeartbeat(String topicName, int partit
315315
heartbeatMonitoringService.getHeartbeatInfo(topicName, partition, false);
316316
leaderSepRTTopicCount += heartbeatInfoMap.keySet().stream().filter(x -> x.endsWith("_sep")).count();
317317
}
318-
Assert.assertEquals(leaderSepRTTopicCount, (long) getNumberOfRegions() * getReplicationFactor());
318+
// Only leaders get heartbeat entries for all regions (including _sep); followers only get local region
319+
Assert.assertEquals(leaderSepRTTopicCount, (long) getNumberOfRegions());
319320
}
320321

321322
private byte[] serializeStringKeyToByteArray(String key) {

0 commit comments

Comments
 (0)