Skip to content

Commit 4cf5954

Browse files
authored
[server][dvc] Propagate pub-sub dependencies to participant consumption state (linkedin#1984)
This change makes pub-sub related dependencies available in ingestion code path, to enable serialization, deserialization, and comparison of PubSubPosition objects.
1 parent c1a37a1 commit 4cf5954

21 files changed

Lines changed: 367 additions & 82 deletions

File tree

clients/da-vinci-client/src/main/java/com/linkedin/davinci/config/VeniceClusterConfig.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@
5454
* class that maintains config very specific to a Venice cluster
5555
*/
5656
public class VeniceClusterConfig {
57-
private static final Logger LOGGER = LogManager.getLogger(VeniceServerConfig.class);
57+
private static final Logger LOGGER = LogManager.getLogger(VeniceClusterConfig.class);
5858

5959
private final String clusterName;
6060
private final String zookeeperAddress;

clients/da-vinci-client/src/main/java/com/linkedin/davinci/config/VeniceServerConfig.java

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,6 @@
4242
import static com.linkedin.venice.ConfigKeys.INGESTION_MLOCK_ENABLED;
4343
import static com.linkedin.venice.ConfigKeys.INGESTION_USE_DA_VINCI_CLIENT;
4444
import static com.linkedin.venice.ConfigKeys.KAFKA_FETCH_THROTTLER_FACTORS_PER_SECOND;
45-
import static com.linkedin.venice.ConfigKeys.KAFKA_PRODUCER_METRICS;
4645
import static com.linkedin.venice.ConfigKeys.KEY_VALUE_PROFILING_ENABLED;
4746
import static com.linkedin.venice.ConfigKeys.KME_REGISTRATION_FROM_MESSAGE_HEADER_ENABLED;
4847
import static com.linkedin.venice.ConfigKeys.LEADER_FOLLOWER_STATE_TRANSITION_THREAD_POOL_STRATEGY;
@@ -489,7 +488,6 @@ public class VeniceServerConfig extends VeniceClusterConfig {
489488
private final long sharedConsumerNonExistingTopicCleanupDelayMS;
490489
private final int offsetLagDeltaRelaxFactorForFastOnlineTransitionInRestart;
491490

492-
private final Set<String> kafkaProducerMetrics;
493491
/**
494492
* Boolean flag indicating if it is a Da Vinci application.
495493
*/
@@ -880,17 +878,6 @@ public VeniceServerConfig(VeniceProperties serverProperties, Map<String, Map<Str
880878
sharedConsumerNonExistingTopicCleanupDelayMS = serverProperties
881879
.getLong(SERVER_SHARED_CONSUMER_NON_EXISTING_TOPIC_CLEANUP_DELAY_MS, TimeUnit.MINUTES.toMillis(10));
882880

883-
List<String> kafkaProducerMetricsList = serverProperties.getList(
884-
KAFKA_PRODUCER_METRICS,
885-
Arrays.asList(
886-
"outgoing-byte-rate",
887-
"record-send-rate",
888-
"batch-size-max",
889-
"batch-size-avg",
890-
"buffer-available-bytes",
891-
"buffer-exhausted-rate"));
892-
kafkaProducerMetrics = new HashSet<>(kafkaProducerMetricsList);
893-
894881
isDaVinciClient = serverProperties.getBoolean(INGESTION_USE_DA_VINCI_CLIENT, false);
895882
unsubscribeAfterBatchpushEnabled = serverProperties.getBoolean(SERVER_UNSUB_AFTER_BATCHPUSH, false);
896883

clients/da-vinci-client/src/main/java/com/linkedin/davinci/kafka/consumer/KafkaStoreIngestionService.java

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,8 @@
5858
import com.linkedin.venice.offsets.OffsetRecord;
5959
import com.linkedin.venice.pubsub.PubSubClientsFactory;
6060
import com.linkedin.venice.pubsub.PubSubConstants;
61+
import com.linkedin.venice.pubsub.PubSubContext;
62+
import com.linkedin.venice.pubsub.PubSubPositionDeserializer;
6163
import com.linkedin.venice.pubsub.PubSubProducerAdapterFactory;
6264
import com.linkedin.venice.pubsub.PubSubTopicPartitionImpl;
6365
import com.linkedin.venice.pubsub.PubSubTopicRepository;
@@ -170,6 +172,7 @@ public class KafkaStoreIngestionService extends AbstractVeniceService implements
170172
private final boolean isIsolatedIngestion;
171173

172174
private final TopicManagerRepository topicManagerRepository;
175+
173176
private ExecutorService participantStoreConsumerExecutorService;
174177

175178
private ExecutorService ingestionExecutorService;
@@ -192,6 +195,7 @@ public class KafkaStoreIngestionService extends AbstractVeniceService implements
192195
private final ResourceAutoClosableLockManager<String> topicLockManager;
193196

194197
private final PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository();
198+
private final PubSubContext pubSubContext;
195199
private final KafkaValueSerializer kafkaValueSerializer;
196200
private final IngestionThrottler ingestionThrottler;
197201
private final ExecutorService aaWCWorkLoadProcessingThreadPool;
@@ -302,6 +306,11 @@ public KafkaStoreIngestionService(
302306
.build();
303307
this.topicManagerRepository =
304308
new TopicManagerRepository(topicManagerContext, serverConfig.getKafkaBootstrapServers());
309+
this.pubSubContext = new PubSubContext.Builder().setTopicManagerRepository(topicManagerRepository)
310+
.setPubSubPositionTypeRegistry(serverConfig.getPubSubPositionTypeRegistry())
311+
.setPubSubPositionDeserializer(new PubSubPositionDeserializer(serverConfig.getPubSubPositionTypeRegistry()))
312+
.setPubSubTopicRepository(pubSubTopicRepository)
313+
.build();
305314

306315
VeniceNotifier notifier = new LogNotifier();
307316
this.leaderFollowerNotifiers.add(notifier);
@@ -487,12 +496,12 @@ public void handleStoreDeleted(Store store) {
487496
serverConfig.getIngestionTaskReusableObjectsStrategy().supplier();
488497

489498
ingestionTaskFactory = StoreIngestionTaskFactory.builder()
499+
.setPubSubContext(pubSubContext)
490500
.setVeniceWriterFactory(veniceWriterFactory)
491501
.setStorageMetadataService(storageMetadataService)
492502
.setLeaderFollowerNotifiersQueue(leaderFollowerNotifiers)
493503
.setSchemaRepository(schemaRepo)
494504
.setMetadataRepository(metadataRepo)
495-
.setTopicManagerRepository(topicManagerRepository)
496505
.setHostLevelIngestionStats(hostLevelIngestionStats)
497506
.setVersionedDIVStats(versionedDIVStats)
498507
.setDaVinciRecordTransformerStats(recordTransformerStats)
@@ -507,7 +516,6 @@ public void handleStoreDeleted(Store store) {
507516
.setMetaStoreWriter(metaStoreWriter)
508517
.setCompressorFactory(compressorFactory)
509518
.setVeniceViewWriterFactory(viewWriterFactory)
510-
.setPubSubTopicRepository(pubSubTopicRepository)
511519
.setRunnableForKillIngestionTasksForNonCurrentVersions(
512520
serverConfig.getIngestionMemoryLimit() > 0 ? () -> killConsumptionTaskForNonCurrentVersions() : null)
513521
.setHeartbeatMonitoringService(heartbeatMonitoringService)
@@ -1407,6 +1415,10 @@ public final ReadOnlyStoreRepository getMetadataRepo() {
14071415
return metadataRepo;
14081416
}
14091417

1418+
public PubSubContext getPubSubContext() {
1419+
return pubSubContext;
1420+
}
1421+
14101422
private boolean ingestionTaskHasAnySubscription(String topic) {
14111423
try (AutoCloseableLock ignore = topicLockManager.getLockForResource(topic)) {
14121424
StoreIngestionTask consumerTask = topicNameToIngestionTaskMap.get(topic);

clients/da-vinci-client/src/main/java/com/linkedin/davinci/kafka/consumer/PartitionConsumptionState.java

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
import com.linkedin.venice.kafka.validation.checksum.CheckSumType;
99
import com.linkedin.venice.meta.Version;
1010
import com.linkedin.venice.offsets.OffsetRecord;
11+
import com.linkedin.venice.pubsub.PubSubContext;
1112
import com.linkedin.venice.pubsub.PubSubTopicPartitionImpl;
1213
import com.linkedin.venice.pubsub.PubSubTopicRepository;
1314
import com.linkedin.venice.pubsub.api.PubSubPosition;
@@ -42,6 +43,7 @@ public class PartitionConsumptionState {
4243
private final int partition;
4344
private final boolean hybrid;
4445
private final OffsetRecord offsetRecord;
46+
private final PubSubContext pubSubContext;
4547

4648
private GUID leaderGUID;
4749

@@ -212,11 +214,17 @@ enum LatchStatus {
212214

213215
private BooleanSupplier isCurrentVersion;
214216

215-
public PartitionConsumptionState(String replicaId, int partition, OffsetRecord offsetRecord, boolean hybrid) {
217+
public PartitionConsumptionState(
218+
String replicaId,
219+
int partition,
220+
OffsetRecord offsetRecord,
221+
PubSubContext pubSubContext,
222+
boolean hybrid) {
216223
this.replicaId = replicaId;
217224
this.partition = partition;
218225
this.hybrid = hybrid;
219226
this.offsetRecord = offsetRecord;
227+
this.pubSubContext = pubSubContext;
220228
this.errorReported = false;
221229
this.lagCaughtUp = false;
222230
this.lagCaughtUpTimeInMs = 0;
@@ -852,4 +860,8 @@ public void clearPendingReportIncPushVersionList() {
852860
pendingReportIncPushVersionList.clear();
853861
offsetRecord.setPendingReportIncPushVersionList(pendingReportIncPushVersionList);
854862
}
863+
864+
public PubSubContext getPubSubContext() {
865+
return pubSubContext;
866+
}
855867
}

clients/da-vinci-client/src/main/java/com/linkedin/davinci/kafka/consumer/StoreIngestionTask.java

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@
9393
import com.linkedin.venice.meta.Store;
9494
import com.linkedin.venice.meta.Version;
9595
import com.linkedin.venice.offsets.OffsetRecord;
96+
import com.linkedin.venice.pubsub.PubSubContext;
9697
import com.linkedin.venice.pubsub.PubSubTopicPartitionImpl;
9798
import com.linkedin.venice.pubsub.PubSubTopicRepository;
9899
import com.linkedin.venice.pubsub.api.DefaultPubSubMessage;
@@ -386,6 +387,7 @@ public abstract class StoreIngestionTask implements Runnable, Closeable {
386387
protected final boolean isChunked;
387388
protected final boolean isRmdChunked;
388389
protected final ChunkedValueManifestSerializer manifestSerializer;
390+
protected final PubSubContext pubSubContext;
389391
protected final PubSubTopicRepository pubSubTopicRepository;
390392
private final String[] msgForLagMeasurement;
391393
private final Runnable runnableForKillIngestionTasksForNonCurrentVersions;
@@ -430,7 +432,9 @@ public StoreIngestionTask(
430432
this.storeRepository = builder.getMetadataRepo();
431433
this.schemaRepository = builder.getSchemaRepo();
432434
this.kafkaVersionTopic = storeVersionConfig.getStoreVersionName();
433-
this.pubSubTopicRepository = builder.getPubSubTopicRepository();
435+
this.pubSubContext = builder.getPubSubContext();
436+
this.pubSubTopicRepository = pubSubContext.getPubSubTopicRepository();
437+
this.topicManagerRepository = pubSubContext.getTopicManagerRepository();
434438
this.versionTopic = pubSubTopicRepository.getTopic(kafkaVersionTopic);
435439
this.storeName = versionTopic.getStoreName();
436440
this.storeVersionName = storeVersionConfig.getStoreVersionName();
@@ -469,7 +473,6 @@ public StoreIngestionTask(
469473
this.consumerDiv = new DataIntegrityValidator(kafkaVersionTopic);
470474
this.consumedBytesSinceLastSync = new VeniceConcurrentHashMap<>();
471475
this.ingestionTaskName = String.format(CONSUMER_TASK_ID_FORMAT, kafkaVersionTopic);
472-
this.topicManagerRepository = builder.getTopicManagerRepository();
473476
this.readOnlyForBatchOnlyStoreEnabled = storeVersionConfig.isReadOnlyForBatchOnlyStoreEnabled();
474477
this.hostLevelIngestionStats = builder.getIngestionStats().getStoreStats(storeName);
475478
this.versionedDIVStats = builder.getVersionedDIVStats();
@@ -2212,6 +2215,7 @@ protected void processCommonConsumerAction(ConsumerAction consumerAction) throws
22122215
getReplicaId(versionTopic, partition),
22132216
partition,
22142217
offsetRecord,
2218+
pubSubContext,
22152219
hybridStoreConfig.isPresent());
22162220
newPartitionConsumptionState.setCurrentVersionSupplier(isCurrentVersion);
22172221

@@ -2375,6 +2379,7 @@ private void resetOffset(int partition, PubSubTopicPartition topicPartition, boo
23752379
getReplicaId(versionTopic, partition),
23762380
partition,
23772381
new OffsetRecord(partitionStateSerializer),
2382+
pubSubContext,
23782383
hybridStoreConfig.isPresent());
23792384
consumptionState.setCurrentVersionSupplier(isCurrentVersion);
23802385
partitionConsumptionStateMap.put(partition, consumptionState);

clients/da-vinci-client/src/main/java/com/linkedin/davinci/kafka/consumer/StoreIngestionTaskFactory.java

Lines changed: 6 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,7 @@
2020
import com.linkedin.venice.meta.ReadOnlyStoreRepository;
2121
import com.linkedin.venice.meta.Store;
2222
import com.linkedin.venice.meta.Version;
23-
import com.linkedin.venice.pubsub.PubSubTopicRepository;
24-
import com.linkedin.venice.pubsub.manager.TopicManagerRepository;
23+
import com.linkedin.venice.pubsub.PubSubContext;
2524
import com.linkedin.venice.serialization.avro.InternalAvroSpecificSerializer;
2625
import com.linkedin.venice.system.store.MetaStoreWriter;
2726
import com.linkedin.venice.utils.DiskUsage;
@@ -111,7 +110,6 @@ public static class Builder {
111110
private Queue<VeniceNotifier> leaderFollowerNotifiers;
112111
private ReadOnlySchemaRepository schemaRepo;
113112
private ReadOnlyStoreRepository metadataRepo;
114-
private TopicManagerRepository topicManagerRepository;
115113
private AggVersionedDaVinciRecordTransformerStats daVinciRecordTransformerStats;
116114
private AggHostLevelIngestionStats ingestionStats;
117115
private AggVersionedDIVStats versionedDIVStats;
@@ -125,7 +123,7 @@ public static class Builder {
125123
private RemoteIngestionRepairService remoteIngestionRepairService;
126124
private MetaStoreWriter metaStoreWriter;
127125
private StorageEngineBackedCompressorFactory compressorFactory;
128-
private PubSubTopicRepository pubSubTopicRepository;
126+
private PubSubContext pubSubContext;
129127
private Runnable runnableForKillIngestionTasksForNonCurrentVersions;
130128
private ExecutorService aaWCWorkLoadProcessingThreadPool;
131129
private ExecutorService aaWCIngestionStorageLookupThreadPool;
@@ -220,12 +218,12 @@ public Builder setMetadataRepository(ReadOnlyStoreRepository metadataRepo) {
220218
return set(() -> this.metadataRepo = metadataRepo);
221219
}
222220

223-
public TopicManagerRepository getTopicManagerRepository() {
224-
return topicManagerRepository;
221+
public Builder setPubSubContext(PubSubContext pubSubContext) {
222+
return set(() -> this.pubSubContext = pubSubContext);
225223
}
226224

227-
public Builder setTopicManagerRepository(TopicManagerRepository topicManagerRepository) {
228-
return set(() -> this.topicManagerRepository = topicManagerRepository);
225+
public PubSubContext getPubSubContext() {
226+
return pubSubContext;
229227
}
230228

231229
public AggVersionedDaVinciRecordTransformerStats getDaVinciRecordTransformerStats() {
@@ -318,14 +316,6 @@ public Builder setCompressorFactory(StorageEngineBackedCompressorFactory compres
318316
return set(() -> this.compressorFactory = compressorFactory);
319317
}
320318

321-
public PubSubTopicRepository getPubSubTopicRepository() {
322-
return pubSubTopicRepository;
323-
}
324-
325-
public Builder setPubSubTopicRepository(PubSubTopicRepository pubSubTopicRepository) {
326-
return set(() -> this.pubSubTopicRepository = pubSubTopicRepository);
327-
}
328-
329319
public Runnable getRunnableForKillIngestionTasksForNonCurrentVersions() {
330320
return runnableForKillIngestionTasksForNonCurrentVersions;
331321
}

clients/da-vinci-client/src/test/java/com/linkedin/davinci/kafka/consumer/ActiveActiveStoreIngestionTaskTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
import com.linkedin.venice.meta.ZKStore;
7070
import com.linkedin.venice.partitioner.DefaultVenicePartitioner;
7171
import com.linkedin.venice.pubsub.ImmutablePubSubMessage;
72+
import com.linkedin.venice.pubsub.PubSubContext;
7273
import com.linkedin.venice.pubsub.PubSubTopicPartitionImpl;
7374
import com.linkedin.venice.pubsub.PubSubTopicRepository;
7475
import com.linkedin.venice.pubsub.adapter.kafka.common.ApacheKafkaOffsetPosition;
@@ -219,7 +220,7 @@ public void testisReadyToServeAnnouncedWithRTLag(IngestionTaskReusableObjects.St
219220

220221
// Set up IngestionTask Builder
221222
StoreIngestionTaskFactory.Builder builder = new StoreIngestionTaskFactory.Builder();
222-
builder.setPubSubTopicRepository(TOPIC_REPOSITORY);
223+
builder.setPubSubContext(new PubSubContext.Builder().setPubSubTopicRepository(TOPIC_REPOSITORY).build());
223224
builder.setHostLevelIngestionStats(mock(AggHostLevelIngestionStats.class));
224225
builder.setAggKafkaConsumerService(mock(AggKafkaConsumerService.class));
225226
builder.setMetadataRepository(readOnlyStoreRepository);

clients/da-vinci-client/src/test/java/com/linkedin/davinci/kafka/consumer/LeaderFollowerStoreIngestionTaskTest.java

Lines changed: 6 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@
5959
import com.linkedin.venice.offsets.InMemoryStorageMetadataService;
6060
import com.linkedin.venice.offsets.OffsetRecord;
6161
import com.linkedin.venice.partitioner.DefaultVenicePartitioner;
62-
import com.linkedin.venice.pubsub.PubSubTopicRepository;
62+
import com.linkedin.venice.pubsub.PubSubContext;
6363
import com.linkedin.venice.pubsub.adapter.kafka.common.ApacheKafkaOffsetPosition;
6464
import com.linkedin.venice.pubsub.api.DefaultPubSubMessage;
6565
import com.linkedin.venice.pubsub.api.PubSubPosition;
@@ -117,6 +117,7 @@ public class LeaderFollowerStoreIngestionTaskTest {
117117
private ReadOnlyStoreRepository storeRepository;
118118
private VeniceServerConfig mockVeniceServerConfig;
119119
private TopicManagerRepository mockTopicManagerRepository;
120+
private PubSubContext pubSubContext;
120121

121122
@Test
122123
public void testCheckWhetherToCloseUnusedVeniceWriter() {
@@ -229,17 +230,12 @@ public void setUp(boolean isHybrid) throws InterruptedException {
229230
.getRefCountedStorageEngine(anyString());
230231
mockVeniceServerConfig = mock(VeniceServerConfig.class);
231232
doReturn(Object2IntMaps.emptyMap()).when(mockVeniceServerConfig).getKafkaClusterUrlToIdMap();
232-
PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository();
233233
hostLevelIngestionStats = mock(HostLevelIngestionStats.class);
234234
AggHostLevelIngestionStats aggHostLevelIngestionStats = mock(AggHostLevelIngestionStats.class);
235235
doReturn(hostLevelIngestionStats).when(aggHostLevelIngestionStats).getStoreStats(storeName);
236236
StorageMetadataService inMemoryStorageMetadataService = new InMemoryStorageMetadataService();
237-
StoreIngestionTaskFactory.Builder builder = getStoreIngestionTaskBuilder(
238-
isHybrid,
239-
storeName,
240-
pubSubTopicRepository,
241-
inMemoryStorageMetadataService,
242-
aggHostLevelIngestionStats);
237+
StoreIngestionTaskFactory.Builder builder =
238+
getStoreIngestionTaskBuilder(isHybrid, storeName, inMemoryStorageMetadataService, aggHostLevelIngestionStats);
243239
when(builder.getSchemaRepo().getKeySchema(storeName)).thenReturn(new SchemaEntry(1, "\"string\""));
244240
mockStore = builder.getMetadataRepo().getStoreOrThrow(storeName);
245241
mockStoreBufferService = (StoreBufferService) builder.getStoreBufferService();
@@ -265,7 +261,8 @@ public void setUp(boolean isHybrid) throws InterruptedException {
265261
doReturn(versionTopic).when(mockVeniceStoreVersionConfig).getStoreVersionName();
266262
mockStorageMetadataService = builder.getStorageMetadataService();
267263
storeRepository = builder.getMetadataRepo();
268-
mockTopicManagerRepository = builder.getTopicManagerRepository();
264+
pubSubContext = builder.getPubSubContext();
265+
mockTopicManagerRepository = pubSubContext.getTopicManagerRepository();
269266
leaderFollowerStoreIngestionTask = spy(
270267
new LeaderFollowerStoreIngestionTask(
271268
mockStorageService,
@@ -287,21 +284,18 @@ public void setUp(boolean isHybrid) throws InterruptedException {
287284
public StoreIngestionTaskFactory.Builder getStoreIngestionTaskBuilder(
288285
boolean isHybrid,
289286
String storeName,
290-
PubSubTopicRepository pubSubTopicRepository,
291287
StorageMetadataService inMemoryStorageMetadataService,
292288
AggHostLevelIngestionStats aggHostLevelIngestionStats) {
293289
if (isHybrid) {
294290
return TestUtils.getStoreIngestionTaskBuilder(storeName, true)
295291
.setServerConfig(mockVeniceServerConfig)
296-
.setPubSubTopicRepository(pubSubTopicRepository)
297292
.setVeniceViewWriterFactory(mockVeniceViewWriterFactory)
298293
.setHeartbeatMonitoringService(mock(HeartbeatMonitoringService.class))
299294
.setCompressorFactory(new StorageEngineBackedCompressorFactory(inMemoryStorageMetadataService))
300295
.setHostLevelIngestionStats(aggHostLevelIngestionStats);
301296
}
302297
return TestUtils.getStoreIngestionTaskBuilder(storeName)
303298
.setServerConfig(mockVeniceServerConfig)
304-
.setPubSubTopicRepository(pubSubTopicRepository)
305299
.setVeniceViewWriterFactory(mockVeniceViewWriterFactory)
306300
.setHeartbeatMonitoringService(mock(HeartbeatMonitoringService.class))
307301
.setCompressorFactory(new StorageEngineBackedCompressorFactory(inMemoryStorageMetadataService))

0 commit comments

Comments
 (0)