Skip to content

Commit b4172e4

Browse files
authored
[server] Remove unused replica progress reporting in ReplicaStatus (linkedin#2100)
Progress was previously reported as a numeric offset. While this was easy to interpret, it was rarely, if ever, used in practice. With the new PubSubPosition, progress is represented as raw bytes, which are not meaningful or helpful for troubleshooting. Instead of adding extra code to report this unused information, the reporting logic is removed.
1 parent 7f17fb5 commit b4172e4

13 files changed

Lines changed: 95 additions & 149 deletions

File tree

clients/da-vinci-client/src/main/java/com/linkedin/davinci/notifier/PushStatusNotifier.java

Lines changed: 12 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
import static com.linkedin.venice.pushmonitor.ExecutionStatus.STARTED;
1010
import static com.linkedin.venice.pushmonitor.ExecutionStatus.START_OF_INCREMENTAL_PUSH_RECEIVED;
1111
import static com.linkedin.venice.pushmonitor.ExecutionStatus.TOPIC_SWITCH_RECEIVED;
12+
import static com.linkedin.venice.pushmonitor.ReplicaStatus.NO_PROGRESS;
1213

1314
import com.linkedin.davinci.config.VeniceServerConfig.IncrementalPushStatusWriteMode;
1415
import com.linkedin.venice.common.PushStatusStoreUtils;
@@ -68,7 +69,7 @@ public void started(String topic, int partitionId, String message) {
6869
@Override
6970
public void restarted(String topic, int partitionId, PubSubPosition position, String message) {
7071
helixPartitionStatusAccessor.updateReplicaStatus(topic, partitionId, STARTED);
71-
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, STARTED, position.getNumericOffset(), "");
72+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, STARTED, "");
7273
}
7374

7475
@Override
@@ -83,7 +84,7 @@ public void completed(String topic, int partitionId, PubSubPosition position, St
8384
LOGGER.error("Could not update CV update to COMPLETED, skipping to update OfflinePushStatus for {}", topic, e);
8485
return;
8586
}
86-
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, COMPLETED, position.getNumericOffset(), "");
87+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, COMPLETED, "");
8788
}
8889

8990
@Override
@@ -101,50 +102,32 @@ public void quotaNotViolated(String topic, int partitionId, PubSubPosition posit
101102
@Override
102103
public void progress(String topic, int partitionId, PubSubPosition position, String message) {
103104
helixPartitionStatusAccessor.updateReplicaStatus(topic, partitionId, PROGRESS);
104-
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, PROGRESS, position.getNumericOffset(), "");
105+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, PROGRESS, "");
105106
}
106107

107108
@Override
108109
public void endOfPushReceived(String topic, int partitionId, PubSubPosition position, String message) {
109-
offLinePushAccessor
110-
.updateReplicaStatus(topic, partitionId, instanceId, END_OF_PUSH_RECEIVED, position.getNumericOffset(), "");
110+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, END_OF_PUSH_RECEIVED, "");
111111
}
112112

113113
@Override
114114
public void topicSwitchReceived(String topic, int partitionId, PubSubPosition position, String message) {
115-
offLinePushAccessor
116-
.updateReplicaStatus(topic, partitionId, instanceId, TOPIC_SWITCH_RECEIVED, position.getNumericOffset(), "");
115+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, TOPIC_SWITCH_RECEIVED, "");
117116
}
118117

119118
@Override
120119
public void dataRecoveryCompleted(String kafkaTopic, int partitionId, PubSubPosition position, String message) {
121-
offLinePushAccessor.updateReplicaStatus(
122-
kafkaTopic,
123-
partitionId,
124-
instanceId,
125-
DATA_RECOVERY_COMPLETED,
126-
position.getNumericOffset(),
127-
message);
120+
offLinePushAccessor.updateReplicaStatus(kafkaTopic, partitionId, instanceId, DATA_RECOVERY_COMPLETED, message);
128121
}
129122

130123
@Override
131124
public void startOfIncrementalPushReceived(String topic, int partitionId, PubSubPosition position, String message) {
132-
updateIncrementalPushStatus(
133-
topic,
134-
partitionId,
135-
position.getNumericOffset(),
136-
message,
137-
START_OF_INCREMENTAL_PUSH_RECEIVED);
125+
updateIncrementalPushStatus(topic, partitionId, NO_PROGRESS, message, START_OF_INCREMENTAL_PUSH_RECEIVED);
138126
}
139127

140128
@Override
141129
public void endOfIncrementalPushReceived(String topic, int partitionId, PubSubPosition position, String message) {
142-
updateIncrementalPushStatus(
143-
topic,
144-
partitionId,
145-
position.getNumericOffset(),
146-
message,
147-
END_OF_INCREMENTAL_PUSH_RECEIVED);
130+
updateIncrementalPushStatus(topic, partitionId, NO_PROGRESS, message, END_OF_INCREMENTAL_PUSH_RECEIVED);
148131
}
149132

150133
private void updateIncrementalPushStatus(
@@ -155,7 +138,7 @@ private void updateIncrementalPushStatus(
155138
ExecutionStatus status) {
156139
if (incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.ZOOKEEPER_ONLY
157140
|| incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.DUAL) {
158-
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, status, offset, message);
141+
offLinePushAccessor.updateReplicaStatus(topic, partitionId, instanceId, status, message);
159142
}
160143
if (incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.PUSH_STATUS_SYSTEM_STORE_ONLY
161144
|| incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.DUAL) {
@@ -172,12 +155,8 @@ public void batchEndOfIncrementalPushReceived(
172155

173156
if (incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.ZOOKEEPER_ONLY
174157
|| incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.DUAL) {
175-
offLinePushAccessor.batchUpdateReplicaIncPushStatus(
176-
topic,
177-
partitionId,
178-
instanceId,
179-
position.getNumericOffset(),
180-
pendingReportIncPushVersionList);
158+
offLinePushAccessor
159+
.batchUpdateReplicaIncPushStatus(topic, partitionId, instanceId, pendingReportIncPushVersionList);
181160
}
182161
if (incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.PUSH_STATUS_SYSTEM_STORE_ONLY
183162
|| incrementalPushStatusWriteMode == IncrementalPushStatusWriteMode.DUAL) {

clients/da-vinci-client/src/test/java/com/linkedin/davinci/notifier/TestPushStatusNotifier.java

Lines changed: 11 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,6 @@
88
import static org.mockito.ArgumentMatchers.any;
99
import static org.mockito.ArgumentMatchers.anyInt;
1010
import static org.mockito.ArgumentMatchers.eq;
11-
import static org.mockito.Mockito.anyLong;
1211
import static org.mockito.Mockito.anyString;
1312
import static org.mockito.Mockito.doReturn;
1413
import static org.mockito.Mockito.doThrow;
@@ -77,12 +76,12 @@ public void testCompleteCVUpdate() {
7776
PubSubPosition p1 = ApacheKafkaOffsetPosition.of(1L);
7877

7978
statusNotifier.completed(topic, 1, p1, "");
80-
verify(offlinePushAccessor, times(1)).updateReplicaStatus(topic, 1, host, ExecutionStatus.COMPLETED, 1, "");
79+
verify(offlinePushAccessor, times(1)).updateReplicaStatus(topic, 1, host, ExecutionStatus.COMPLETED, "");
8180

8281
doThrow(HelixException.class).when(helixPartitionStatusAccessor)
8382
.updateReplicaStatus(any(), anyInt(), eq(ExecutionStatus.COMPLETED));
8483
statusNotifier.completed(topic, 1, p1, "");
85-
verify(offlinePushAccessor, never()).updateReplicaStatus(topic, 1, "host", ExecutionStatus.COMPLETED, 1, "");
84+
verify(offlinePushAccessor, never()).updateReplicaStatus(topic, 1, "host", ExecutionStatus.COMPLETED, "");
8685

8786
doReturn(mock(Store.class)).when(storeRepository).getStoreOrThrow(any());
8887
statusNotifier.startOfIncrementalPushReceived(topic, 1, p1, "");
@@ -119,16 +118,11 @@ public void testStartOfIncrementalPushReceived(
119118
notifier.startOfIncrementalPushReceived(TOPIC, PARTITION_ID, POSITION_12345, MESSAGE);
120119

121120
if (expectZookeeper) {
122-
verify(offlinePushAccessor, times(1)).updateReplicaStatus(
123-
TOPIC,
124-
PARTITION_ID,
125-
INSTANCE_ID,
126-
START_OF_INCREMENTAL_PUSH_RECEIVED,
127-
POSITION_12345.getInternalOffset(),
128-
MESSAGE);
121+
verify(offlinePushAccessor, times(1))
122+
.updateReplicaStatus(TOPIC, PARTITION_ID, INSTANCE_ID, START_OF_INCREMENTAL_PUSH_RECEIVED, MESSAGE);
129123
} else {
130124
verify(offlinePushAccessor, never())
131-
.updateReplicaStatus(anyString(), anyInt(), anyString(), any(ExecutionStatus.class), anyLong(), anyString());
125+
.updateReplicaStatus(anyString(), anyInt(), anyString(), any(ExecutionStatus.class), anyString());
132126
}
133127

134128
if (expectPushStatusStore) {
@@ -162,16 +156,11 @@ public void testEndOfIncrementalPushReceived(
162156
notifier.endOfIncrementalPushReceived(TOPIC, PARTITION_ID, POSITION_12345, MESSAGE);
163157

164158
if (expectZookeeper) {
165-
verify(offlinePushAccessor, times(1)).updateReplicaStatus(
166-
TOPIC,
167-
PARTITION_ID,
168-
INSTANCE_ID,
169-
END_OF_INCREMENTAL_PUSH_RECEIVED,
170-
POSITION_12345.getInternalOffset(),
171-
MESSAGE);
159+
verify(offlinePushAccessor, times(1))
160+
.updateReplicaStatus(TOPIC, PARTITION_ID, INSTANCE_ID, END_OF_INCREMENTAL_PUSH_RECEIVED, MESSAGE);
172161
} else {
173162
verify(offlinePushAccessor, never())
174-
.updateReplicaStatus(anyString(), anyInt(), anyString(), any(ExecutionStatus.class), anyLong(), anyString());
163+
.updateReplicaStatus(anyString(), anyInt(), anyString(), any(ExecutionStatus.class), anyString());
175164
}
176165

177166
if (expectPushStatusStore) {
@@ -211,15 +200,10 @@ public void testBatchEndOfIncrementalPushReceived(
211200
notifier.batchEndOfIncrementalPushReceived(TOPIC, PARTITION_ID, POSITION_12345, incPushVersions);
212201

213202
if (expectZookeeper) {
214-
verify(offlinePushAccessor, times(1)).batchUpdateReplicaIncPushStatus(
215-
TOPIC,
216-
PARTITION_ID,
217-
INSTANCE_ID,
218-
POSITION_12345.getInternalOffset(),
219-
incPushVersions);
203+
verify(offlinePushAccessor, times(1))
204+
.batchUpdateReplicaIncPushStatus(TOPIC, PARTITION_ID, INSTANCE_ID, incPushVersions);
220205
} else {
221-
verify(offlinePushAccessor, never())
222-
.batchUpdateReplicaIncPushStatus(anyString(), anyInt(), anyString(), anyLong(), any());
206+
verify(offlinePushAccessor, never()).batchUpdateReplicaIncPushStatus(anyString(), anyInt(), anyString(), any());
223207
}
224208

225209
if (expectPushStatusStore) {

internal/venice-common/src/main/java/com/linkedin/venice/helix/VeniceOfflinePushMonitorAccessor.java

Lines changed: 4 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@
1010
import com.linkedin.venice.pushmonitor.PartitionStatus;
1111
import com.linkedin.venice.pushmonitor.PartitionStatusListener;
1212
import com.linkedin.venice.pushmonitor.ReadOnlyPartitionStatus;
13-
import com.linkedin.venice.pushmonitor.ReplicaStatus;
1413
import com.linkedin.venice.utils.HelixUtils;
1514
import com.linkedin.venice.utils.LogContext;
1615
import com.linkedin.venice.utils.PathResourceRegistry;
@@ -227,35 +226,17 @@ public void updateReplicaStatus(
227226
int partitionId,
228227
String instanceId,
229228
ExecutionStatus status,
230-
long progress,
231229
String incrementalPushVersion) {
232-
compareAndUpdateReplicaStatus(topic, partitionId, instanceId, status, progress, incrementalPushVersion);
233-
}
234-
235-
@Override
236-
public void updateReplicaStatus(
237-
String topic,
238-
int partitionId,
239-
String instanceId,
240-
ExecutionStatus status,
241-
String incrementalPushVersion) {
242-
compareAndUpdateReplicaStatus(
243-
topic,
244-
partitionId,
245-
instanceId,
246-
status,
247-
ReplicaStatus.NO_PROGRESS,
248-
incrementalPushVersion);
230+
compareAndUpdateReplicaStatus(topic, partitionId, instanceId, status, incrementalPushVersion);
249231
}
250232

251233
@Override
252234
public void batchUpdateReplicaIncPushStatus(
253235
String kafkaTopic,
254236
int partitionId,
255237
String instanceId,
256-
long progress,
257238
List<String> pendingReportIncPushVersionList) {
258-
compareAndBatchUpdateReplicaStatus(kafkaTopic, partitionId, instanceId, progress, pendingReportIncPushVersionList);
239+
compareAndBatchUpdateReplicaStatus(kafkaTopic, partitionId, instanceId, pendingReportIncPushVersionList);
259240
}
260241

261242
/**
@@ -274,7 +255,6 @@ private void compareAndUpdateReplicaStatus(
274255
int partitionId,
275256
String instanceId,
276257
ExecutionStatus status,
277-
long progress,
278258
String incrementalPushVersion) {
279259
// If a version was created prior to the deployment of this new push monitor, an exception would be thrown while
280260
// upgrading venice server.
@@ -296,7 +276,7 @@ private void compareAndUpdateReplicaStatus(
296276
currentData = new PartitionStatus(partitionId);
297277
}
298278

299-
currentData.updateReplicaStatus(instanceId, status, incrementalPushVersion, progress);
279+
currentData.updateReplicaStatus(instanceId, status, incrementalPushVersion);
300280
return currentData;
301281
});
302282
LOGGER.info("Updated replica status for replica: {} status: {} in cluster: {}.", replicaId, status, clusterName);
@@ -306,7 +286,6 @@ private void compareAndBatchUpdateReplicaStatus(
306286
String topic,
307287
int partitionId,
308288
String instanceId,
309-
long progress,
310289
List<String> incPushBatchStatus) {
311290
if (!pushStatusExists(topic)) {
312291
return;
@@ -317,7 +296,7 @@ private void compareAndBatchUpdateReplicaStatus(
317296
if (currentData == null) {
318297
currentData = new PartitionStatus(partitionId);
319298
}
320-
currentData.batchUpdateReplicaIncPushStatus(instanceId, incPushBatchStatus, progress);
299+
currentData.batchUpdateReplicaIncPushStatus(instanceId, incPushBatchStatus);
321300
return currentData;
322301
});
323302
LOGGER.info(

internal/venice-common/src/main/java/com/linkedin/venice/pushmonitor/OfflinePushAccessor.java

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -49,18 +49,7 @@ public interface OfflinePushAccessor {
4949
void deleteOfflinePushStatusAndItsPartitionStatuses(String kafkaTopic);
5050

5151
/**
52-
* Update one particular replica status and progress by given topic, partition and instanceId to the persistent storage.
53-
*/
54-
void updateReplicaStatus(
55-
String kafkaTopic,
56-
int partitionId,
57-
String instanceId,
58-
ExecutionStatus status,
59-
long progress,
60-
String message);
61-
62-
/**
63-
* Update one particular replica status only by given topic, partition and instanceId to the persistent storage.
52+
* Update one particular replica status by given topic, partition and instanceId to the persistent storage.
6453
*/
6554
void updateReplicaStatus(
6655
String kafkaTopic,
@@ -73,7 +62,6 @@ default void batchUpdateReplicaIncPushStatus(
7362
String kafkaTopic,
7463
int partitionId,
7564
String instanceId,
76-
long progress,
7765
List<String> pendingReportIncPushVersionList) {
7866
}
7967

internal/venice-common/src/main/java/com/linkedin/venice/pushmonitor/PartitionStatus.java

Lines changed: 4 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -38,26 +38,13 @@ public void updateReplicaStatus(String instanceId, ExecutionStatus newStatus, bo
3838
updateReplicaStatus(instanceId, newStatus, "", enableStatusHistory);
3939
}
4040

41-
public void updateReplicaStatus(
42-
String instanceId,
43-
ExecutionStatus newStatus,
44-
String incrementalPushVersion,
45-
long progress) {
46-
ReplicaStatus replicaStatus = updateReplicaStatus(instanceId, newStatus, incrementalPushVersion, true);
47-
replicaStatus.setCurrentProgress(progress);
41+
public void updateReplicaStatus(String instanceId, ExecutionStatus newStatus, String incrementalPushVersion) {
42+
updateReplicaStatus(instanceId, newStatus, incrementalPushVersion, true);
4843
}
4944

50-
public void batchUpdateReplicaIncPushStatus(String instanceId, List<String> incPushVersionList, long progress) {
51-
ReplicaStatus replicaStatus = null;
45+
public void batchUpdateReplicaIncPushStatus(String instanceId, List<String> incPushVersionList) {
5246
for (String incrementalPushVersion: incPushVersionList) {
53-
replicaStatus = updateReplicaStatus(
54-
instanceId,
55-
ExecutionStatus.END_OF_INCREMENTAL_PUSH_RECEIVED,
56-
incrementalPushVersion,
57-
true);
58-
}
59-
if (replicaStatus != null) {
60-
replicaStatus.setCurrentProgress(progress);
47+
updateReplicaStatus(instanceId, ExecutionStatus.END_OF_INCREMENTAL_PUSH_RECEIVED, incrementalPushVersion, true);
6148
}
6249
}
6350

internal/venice-common/src/main/java/com/linkedin/venice/pushmonitor/ReadOnlyPartitionStatus.java

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -17,16 +17,17 @@ public void updateReplicaStatus(String instanceId, ExecutionStatus newStatus) {
1717

1818
@Override
1919
public void updateReplicaStatus(String instanceId, ExecutionStatus newStatus, boolean enableStatusHistory) {
20-
throw new VeniceException("Unsupported operation in ReadonlyPartition status: updateProgress.");
20+
throw new VeniceException("Unsupported operation in ReadonlyPartition status: updateReplicaStatus.");
21+
}
22+
23+
@Override
24+
public void updateReplicaStatus(String instanceId, ExecutionStatus newStatus, String incrementalPushVersion) {
25+
throw new VeniceException("Unsupported operation in ReadonlyPartition status: updateReplicaStatus.");
2126
}
2227

2328
@Override
24-
public void updateReplicaStatus(
25-
String instanceId,
26-
ExecutionStatus newStatus,
27-
String incrementalPushVersion,
28-
long progress) {
29-
throw new VeniceException("Unsupported operation in ReadonlyPartition status: updateProgress.");
29+
public void batchUpdateReplicaIncPushStatus(String instanceId, java.util.List<String> incPushVersionList) {
30+
throw new VeniceException("Unsupported operation in ReadonlyPartition status: batchUpdateReplicaIncPushStatus.");
3031
}
3132

3233
public static ReadOnlyPartitionStatus fromPartitionStatus(PartitionStatus partitionStatus) {

0 commit comments

Comments
 (0)