Skip to content

Commit ba39e96

Browse files
authored
[server][dvc][cc] Read checkpoint state from PubSubPosition field with offset fallback (linkedin#2123)
Dual-write to both offset and PubSubPosition fields was already in place. This change switches checkpoint reads in OffsetRecord to use the PubSubPosition field by default, with a fallback to the legacy offset field for backward compatibility.
1 parent 7398e12 commit ba39e96

49 files changed

Lines changed: 586 additions & 274 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

clients/da-vinci-client/src/main/java/com/linkedin/davinci/blobtransfer/BlobSnapshotManager.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -342,7 +342,7 @@ public BlobTransferPartitionMetadata prepareMetadata(BlobTransferPayload blobTra
342342

343343
if (storageMetadataService.getStoreVersionState(blobTransferRequest.getTopicName()) == null
344344
|| storageMetadataService
345-
.getLastOffset(blobTransferRequest.getTopicName(), blobTransferRequest.getPartition()) == null) {
345+
.getLastOffset(blobTransferRequest.getTopicName(), blobTransferRequest.getPartition(), null) == null) {
346346
throw new VeniceException("Cannot get store version state or offset record from storage metadata service.");
347347
}
348348

@@ -352,8 +352,8 @@ public BlobTransferPartitionMetadata prepareMetadata(BlobTransferPayload blobTra
352352
java.nio.ByteBuffer storeVersionStateByte =
353353
ByteBuffer.wrap(storeVersionStateSerializer.serialize(blobTransferRequest.getTopicName(), storeVersionState));
354354

355-
OffsetRecord offsetRecord =
356-
storageMetadataService.getLastOffset(blobTransferRequest.getTopicName(), blobTransferRequest.getPartition());
355+
OffsetRecord offsetRecord = storageMetadataService
356+
.getLastOffset(blobTransferRequest.getTopicName(), blobTransferRequest.getPartition(), null);
357357
java.nio.ByteBuffer offsetRecordByte = ByteBuffer.wrap(offsetRecord.toBytes());
358358

359359
return new BlobTransferPartitionMetadata(

clients/da-vinci-client/src/main/java/com/linkedin/davinci/blobtransfer/client/P2PMetadataTransferHandler.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -107,7 +107,7 @@ public void updateStorePartitionMetadata(
107107
storageMetadataService.put(
108108
transferredPartitionMetadata.topicName,
109109
transferredPartitionMetadata.partitionId,
110-
new OffsetRecord(transferredPartitionMetadata.offsetRecord.array(), partitionStateSerializer));
110+
new OffsetRecord(transferredPartitionMetadata.offsetRecord.array(), partitionStateSerializer, null));
111111
// update the metadata SVS
112112
updateStorageVersionState(storageMetadataService, transferredPartitionMetadata);
113113
}

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

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import com.linkedin.venice.compression.VeniceCompressor;
77
import com.linkedin.venice.exceptions.VeniceException;
88
import com.linkedin.venice.kafka.protocol.state.PartitionState;
9+
import com.linkedin.venice.pubsub.PubSubContext;
910
import com.linkedin.venice.serialization.avro.InternalAvroSpecificSerializer;
1011
import com.linkedin.venice.utils.lazy.Lazy;
1112
import java.io.Closeable;
@@ -312,8 +313,10 @@ public final void onRecovery(
312313
StorageEngine storageEngine,
313314
int partitionId,
314315
InternalAvroSpecificSerializer<PartitionState> partitionStateSerializer,
315-
Lazy<VeniceCompressor> compressor) {
316-
recordTransformerUtility.onRecovery(storageEngine, partitionId, partitionStateSerializer, compressor);
316+
Lazy<VeniceCompressor> compressor,
317+
PubSubContext pubSubContext) {
318+
recordTransformerUtility
319+
.onRecovery(storageEngine, partitionId, partitionStateSerializer, compressor, pubSubContext);
317320
}
318321

319322
/**

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

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import com.linkedin.venice.compression.VeniceCompressor;
88
import com.linkedin.venice.kafka.protocol.state.PartitionState;
99
import com.linkedin.venice.offsets.OffsetRecord;
10+
import com.linkedin.venice.pubsub.PubSubContext;
1011
import com.linkedin.venice.serialization.avro.InternalAvroSpecificSerializer;
1112
import com.linkedin.venice.serializer.FastSerializerDeserializerFactory;
1213
import com.linkedin.venice.serializer.RecordDeserializer;
@@ -128,10 +129,12 @@ public final void onRecovery(
128129
StorageEngine storageEngine,
129130
int partitionId,
130131
InternalAvroSpecificSerializer<PartitionState> partitionStateSerializer,
131-
Lazy<VeniceCompressor> compressor) {
132+
Lazy<VeniceCompressor> compressor,
133+
PubSubContext pubSubContext) {
132134
int classHash = recordTransformer.getClassHash();
133-
Optional<OffsetRecord> optionalOffsetRecord = storageEngine.getPartitionOffset(partitionId);
134-
OffsetRecord offsetRecord = optionalOffsetRecord.orElseGet(() -> new OffsetRecord(partitionStateSerializer));
135+
Optional<OffsetRecord> optionalOffsetRecord = storageEngine.getPartitionOffset(partitionId, pubSubContext);
136+
OffsetRecord offsetRecord =
137+
optionalOffsetRecord.orElseGet(() -> new OffsetRecord(partitionStateSerializer, pubSubContext));
135138

136139
boolean transformerLogicChanged = hasTransformerLogicChanged(classHash, offsetRecord);
137140

@@ -142,7 +145,7 @@ public final void onRecovery(
142145
storageEngine.clearPartitionOffset(partitionId);
143146

144147
// Offset record is deleted, so create a new one and persist it
145-
offsetRecord = new OffsetRecord(partitionStateSerializer);
148+
offsetRecord = new OffsetRecord(partitionStateSerializer, pubSubContext);
146149
offsetRecord.setRecordTransformerClassHash(classHash);
147150
storageEngine.putPartitionOffset(partitionId, offsetRecord);
148151
} else {

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

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import com.linkedin.venice.annotation.Experimental;
66
import com.linkedin.venice.compression.VeniceCompressor;
77
import com.linkedin.venice.kafka.protocol.state.PartitionState;
8+
import com.linkedin.venice.pubsub.PubSubContext;
89
import com.linkedin.venice.serialization.avro.InternalAvroSpecificSerializer;
910
import com.linkedin.venice.utils.lazy.Lazy;
1011
import java.io.IOException;
@@ -110,8 +111,9 @@ public void internalOnRecovery(
110111
StorageEngine storageEngine,
111112
int partitionId,
112113
InternalAvroSpecificSerializer<PartitionState> partitionStateSerializer,
113-
Lazy<VeniceCompressor> compressor) {
114-
this.recordTransformer.onRecovery(storageEngine, partitionId, partitionStateSerializer, compressor);
114+
Lazy<VeniceCompressor> compressor,
115+
PubSubContext pubSubContext) {
116+
this.recordTransformer.onRecovery(storageEngine, partitionId, partitionStateSerializer, compressor, pubSubContext);
115117
}
116118

117119
public long getCountDownStartConsumptionLatchCount() {

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

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import com.linkedin.venice.controllerapi.D2ControllerClient;
66
import com.linkedin.venice.pubsub.PubSubClientsFactory;
77
import com.linkedin.venice.pubsub.PubSubConsumerAdapterFactory;
8+
import com.linkedin.venice.pubsub.PubSubContext;
89
import com.linkedin.venice.pubsub.PubSubPositionDeserializer;
910
import com.linkedin.venice.pubsub.PubSubPositionTypeRegistry;
1011
import com.linkedin.venice.pubsub.PubSubTopicRepository;
@@ -67,11 +68,8 @@ public class ChangelogClientConfig<T extends SpecificRecord> {
6768
* These are refreshed each time a new set of consumer properties is applied.
6869
*/
6970
private PubSubConsumerAdapterFactory<? extends PubSubConsumerAdapter> pubSubConsumerAdapterFactory;
70-
private PubSubPositionDeserializer pubSubPositionDeserializer;
7171
private PubSubMessageDeserializer pubSubMessageDeserializer;
72-
private PubSubPositionTypeRegistry pubSubPositionTypeRegistry;
73-
74-
private final PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository();
72+
private PubSubContext pubSubContext;
7573

7674
public ChangelogClientConfig(String storeName) {
7775
this.innerClientConfig = new ClientConfig<>(storeName);
@@ -358,10 +356,10 @@ public ChangelogClientConfig setMaxBufferSize(int maxBufferSize) {
358356
*
359357
* <p>This method sets up:
360358
* <ul>
361-
* <li>{@link #pubSubPositionTypeRegistry} – derived from consumer properties</li>
362-
* <li>{@link #pubSubPositionDeserializer} – uses the initialized position type registry</li>
363-
* <li>{@link #pubSubConsumerAdapterFactory} – created based on resolved consumer configuration</li>
364-
* <li>{@link #pubSubMessageDeserializer} – stateless shared instance</li>
359+
* <li>{@link PubSubPositionTypeRegistry} – derived from consumer properties</li>
360+
* <li>{@link PubSubPositionDeserializer} – uses the initialized position type registry</li>
361+
* <li>{@link PubSubConsumerAdapterFactory} – created based on resolved consumer configuration</li>
362+
* <li>{@link PubSubMessageDeserializer} – stateless shared instance</li>
365363
* </ul>
366364
*
367365
* <p><strong>Note:</strong> These fields are derived from the {@link #consumerProperties} and should
@@ -371,8 +369,14 @@ public ChangelogClientConfig setMaxBufferSize(int maxBufferSize) {
371369
*/
372370
private void initializePubSubInternals() {
373371
VeniceProperties pubSubProperties = new VeniceProperties(this.consumerProperties);
374-
this.pubSubPositionTypeRegistry = PubSubPositionTypeRegistry.fromPropertiesOrDefault(pubSubProperties);
375-
this.pubSubPositionDeserializer = new PubSubPositionDeserializer(pubSubPositionTypeRegistry);
372+
PubSubPositionTypeRegistry typeRegistry = PubSubPositionTypeRegistry.fromPropertiesOrDefault(pubSubProperties);
373+
PubSubPositionDeserializer pubSubPositionDeserializer = new PubSubPositionDeserializer(typeRegistry);
374+
PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository();
375+
// todo(sushantmane): Consider passing TopicManagerRepository from outside if required.
376+
this.pubSubContext = new PubSubContext.Builder().setPubSubPositionDeserializer(pubSubPositionDeserializer)
377+
.setPubSubPositionTypeRegistry(typeRegistry)
378+
.setPubSubTopicRepository(pubSubTopicRepository)
379+
.build();
376380
this.pubSubConsumerAdapterFactory = PubSubClientsFactory.createConsumerFactory(pubSubProperties);
377381
this.pubSubMessageDeserializer = PubSubMessageDeserializer.createOptimizedDeserializer();
378382
}
@@ -381,20 +385,12 @@ protected PubSubConsumerAdapterFactory<? extends PubSubConsumerAdapter> getPubSu
381385
return pubSubConsumerAdapterFactory;
382386
}
383387

384-
protected PubSubPositionDeserializer getPubSubPositionDeserializer() {
385-
return pubSubPositionDeserializer;
386-
}
387-
388388
protected PubSubMessageDeserializer getPubSubMessageDeserializer() {
389389
return pubSubMessageDeserializer;
390390
}
391391

392-
protected PubSubPositionTypeRegistry getPubSubPositionTypeRegistry() {
393-
return pubSubPositionTypeRegistry;
394-
}
395-
396-
protected PubSubTopicRepository getPubSubTopicRepository() {
397-
return pubSubTopicRepository;
392+
protected PubSubContext getPubSubContext() {
393+
return pubSubContext;
398394
}
399395

400396
private ChangelogClientConfig setInnerClientConfig(ClientConfig<T> innerClientConfig) {

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -150,7 +150,7 @@ private Function<String, Boolean> functionToCheckWhetherStorageEngineShouldBeKep
150150
// subscribe to that position. If it's not able to, that means that the local state is off Venice retention,
151151
// and therefore should be completely re-bootstrapped.
152152
for (Integer partition: bootstrapStateMap.keySet()) {
153-
OffsetRecord offsetRecord = storageMetadataService.getLastOffset(localStateTopicName, partition);
153+
OffsetRecord offsetRecord = storageMetadataService.getLastOffset(localStateTopicName, partition, pubSubContext);
154154
if (offsetRecord == null) {
155155
// No offset info in local, need to bootstrap from beginning.
156156
return false;
@@ -304,7 +304,7 @@ public void onCompletion() {
304304
* This method flushes data partition on disk and syncs the underlying database with {@link OffsetRecord}.
305305
*/
306306
private void syncOffset(int partitionId, BootstrapState bootstrapState) {
307-
OffsetRecord lastOffset = storageMetadataService.getLastOffset(localStateTopicName, partitionId);
307+
OffsetRecord lastOffset = storageMetadataService.getLastOffset(localStateTopicName, partitionId, pubSubContext);
308308
StorageEngine storageEngineReloadedFromRepo = getStorageEngine(localStateTopicName);
309309
if (storageEngineReloadedFromRepo == null) {
310310
String replicaId = Utils.getReplicaId(localStateTopicName, partitionId);
@@ -486,7 +486,7 @@ public CompletableFuture<Void> seekWithBootStrap(Set<Integer> partitions) {
486486
partition,
487487
() -> null);
488488
// Get the last persisted Offset record from metadata service
489-
OffsetRecord offsetRecord = storageMetadataService.getLastOffset(localStateTopicName, partition);
489+
OffsetRecord offsetRecord = storageMetadataService.getLastOffset(localStateTopicName, partition, pubSubContext);
490490
// Where we're at now
491491
String offsetString = offsetRecord.getDatabaseInfo().get(CHANGE_CAPTURE_COORDINATE);
492492
VeniceChangeCoordinate localCheckpoint;

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -219,8 +219,8 @@ protected static PubSubConsumerAdapter getPubSubConsumer(
219219
PubSubConsumerAdapterContext context = new PubSubConsumerAdapterContext.Builder().setConsumerName(consumerName)
220220
.setVeniceProperties(new VeniceProperties(changelogClientConfig.getConsumerProperties()))
221221
.setPubSubMessageDeserializer(changelogClientConfig.getPubSubMessageDeserializer())
222-
.setPubSubTopicRepository(changelogClientConfig.getPubSubTopicRepository())
223-
.setPubSubPositionTypeRegistry(changelogClientConfig.getPubSubPositionTypeRegistry())
222+
.setPubSubTopicRepository(changelogClientConfig.getPubSubContext().getPubSubTopicRepository())
223+
.setPubSubPositionTypeRegistry(changelogClientConfig.getPubSubContext().getPubSubPositionTypeRegistry())
224224
.build();
225225
return changelogClientConfig.getPubSubConsumerAdapterFactory().create(context);
226226
}

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

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@
4646
import com.linkedin.venice.meta.PersistenceType;
4747
import com.linkedin.venice.meta.Store;
4848
import com.linkedin.venice.meta.Version;
49+
import com.linkedin.venice.pubsub.PubSubContext;
4950
import com.linkedin.venice.pubsub.PubSubPositionDeserializer;
5051
import com.linkedin.venice.pubsub.PubSubTopicPartitionImpl;
5152
import com.linkedin.venice.pubsub.PubSubTopicRepository;
@@ -150,6 +151,7 @@ public class VeniceChangelogConsumerImpl<K, V> implements VeniceChangelogConsume
150151
protected final PubSubConsumerAdapter pubSubConsumer;
151152
protected final PubSubTopicRepository pubSubTopicRepository;
152153
protected final PubSubPositionDeserializer pubSubPositionDeserializer;
154+
protected final PubSubContext pubSubContext;
153155
protected final ExecutorService seekExecutorService;
154156

155157
// This member is a map of maps in order to accommodate view topics. If the message we consume has the appropriate
@@ -181,8 +183,9 @@ public VeniceChangelogConsumerImpl(
181183
long consumerSequenceIdStartingValue) {
182184
Objects.requireNonNull(changelogClientConfig, "ChangelogClientConfig cannot be null");
183185
this.pubSubConsumer = pubSubConsumer;
184-
this.pubSubTopicRepository = changelogClientConfig.getPubSubTopicRepository();
185-
this.pubSubPositionDeserializer = changelogClientConfig.getPubSubPositionDeserializer();
186+
this.pubSubContext = changelogClientConfig.getPubSubContext();
187+
this.pubSubTopicRepository = pubSubContext.getPubSubTopicRepository();
188+
this.pubSubPositionDeserializer = pubSubContext.getPubSubPositionDeserializer();
186189

187190
seekExecutorService = Executors.newFixedThreadPool(10);
188191

clients/da-vinci-client/src/main/java/com/linkedin/davinci/ingestion/DefaultIngestionBackend.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -326,7 +326,8 @@ public boolean isOffsetLagged(
326326
long blobTransferDisabledOffsetLagThreshold,
327327
boolean hybridStore) {
328328
String topicName = Version.composeKafkaTopic(store, versionNumber);
329-
OffsetRecord offsetRecord = storageMetadataService.getLastOffset(topicName, partition);
329+
OffsetRecord offsetRecord =
330+
storageMetadataService.getLastOffset(topicName, partition, storeIngestionService.getPubSubContext());
330331

331332
if (offsetRecord == null) {
332333
return true;

0 commit comments

Comments
 (0)