Skip to content

Commit deaae9e

Browse files
authored
[protocol][fast-client][server][build][compat] External-storage hooks + MetadataResponseRecord v3 activation (linkedin#2830)
1 parent 3367b8d commit deaae9e

23 files changed

Lines changed: 1592 additions & 16 deletions

File tree

build.gradle

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -331,11 +331,7 @@ subprojects {
331331
def versionOverrides = [
332332
// AdminOperation v100 stages degradedDatacenters on AddVersion.
333333
// Pinned to the current active version until the Java wiring lands in a follow-up PR.
334-
project(':services:venice-controller').file('src/main/resources/avro/AdminOperation/v99', PathValidation.DIRECTORY),
335-
// MetadataResponseRecord v3 stages externalStorageReadMode (mirrors StoreProperties.externalStorageReadMode
336-
// from StoreMetaValue v44). Pinned to the current active version until SERVER_METADATA_RESPONSE is bumped
337-
// and server populate + Fast Client read land in follow-up PRs.
338-
project(':internal:venice-common').file('src/main/resources/avro/MetadataResponseRecord/v2', PathValidation.DIRECTORY)
334+
project(':services:venice-controller').file('src/main/resources/avro/AdminOperation/v99', PathValidation.DIRECTORY)
339335
]
340336

341337
def schemaDirs = [sourceDir]

clients/da-vinci-client/src/main/java/com/linkedin/davinci/listener/response/MetadataResponse.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,10 @@ public void setBatchGetLimit(int batchGetLimit) {
5656
responseRecord.setBatchGetLimit(batchGetLimit);
5757
}
5858

59+
public void setExternalStorageReadMode(int externalStorageReadMode) {
60+
responseRecord.setExternalStorageReadMode(externalStorageReadMode);
61+
}
62+
5963
public ByteBuf getResponseBody() {
6064
return Unpooled.wrappedBuffer(serializedResponse());
6165
}

clients/venice-client/src/main/java/com/linkedin/venice/fastclient/DelegatingAvroStoreClient.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,10 +6,13 @@
66
import com.linkedin.venice.client.store.AvroGenericStoreClient;
77
import com.linkedin.venice.client.store.ComputeGenericRecord;
88
import com.linkedin.venice.client.store.ComputeRequestBuilder;
9+
import com.linkedin.venice.client.store.listeners.StoreConfigChangeListener;
10+
import com.linkedin.venice.client.store.listeners.StoreVersionSwitchListener;
911
import com.linkedin.venice.client.store.streaming.StreamingCallback;
1012
import com.linkedin.venice.compute.ComputeRequestWrapper;
1113
import com.linkedin.venice.fastclient.factory.ClientFactory;
1214
import com.linkedin.venice.schema.SchemaReader;
15+
import java.nio.ByteBuffer;
1316
import java.util.Optional;
1417
import java.util.Set;
1518
import java.util.concurrent.CompletableFuture;
@@ -144,4 +147,19 @@ public ComputeRequestBuilder<K> compute(
144147
long preRequestTimeInNS) throws VeniceClientException {
145148
return delegate.compute(stats, streamingStats, computeStoreClient, preRequestTimeInNS);
146149
}
150+
151+
@Override
152+
public V decompressAndDeserialize(ByteBuffer rawValue, int version, K key) throws VeniceClientException {
153+
return delegate.decompressAndDeserialize(rawValue, version, key);
154+
}
155+
156+
@Override
157+
public void registerVersionSwitchListener(StoreVersionSwitchListener listener) {
158+
delegate.registerVersionSwitchListener(listener);
159+
}
160+
161+
@Override
162+
public void registerStoreConfigChangeListener(StoreConfigChangeListener listener) {
163+
delegate.registerStoreConfigChangeListener(listener);
164+
}
147165
}

clients/venice-client/src/main/java/com/linkedin/venice/fastclient/DispatchingAvroGenericStoreClient.java

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,8 @@
1212
import com.linkedin.venice.client.stats.ClientStats;
1313
import com.linkedin.venice.client.store.AbstractAvroStoreClient;
1414
import com.linkedin.venice.client.store.ComputeGenericRecord;
15+
import com.linkedin.venice.client.store.listeners.StoreConfigChangeListener;
16+
import com.linkedin.venice.client.store.listeners.StoreVersionSwitchListener;
1517
import com.linkedin.venice.client.store.streaming.ComputeRecordStreamDecoder;
1618
import com.linkedin.venice.client.store.streaming.StreamingCallback;
1719
import com.linkedin.venice.client.store.streaming.TrackingStreamingCallback;
@@ -41,6 +43,7 @@
4143
import com.linkedin.venice.utils.LatencyUtils;
4244
import com.linkedin.venice.utils.RedundantExceptionFilter;
4345
import com.linkedin.venice.utils.concurrent.ChainedCompletableFuture;
46+
import java.io.IOException;
4447
import java.nio.ByteBuffer;
4548
import java.util.ArrayList;
4649
import java.util.Collections;
@@ -792,4 +795,45 @@ public Schema getLatestValueSchema() {
792795
public SchemaReader getSchemaReader() {
793796
return metadata;
794797
}
798+
799+
/**
800+
* Fast Client implementation of the external-storage re-entry seam. Reads the 4-byte BE writer-schema-id prefix
801+
* from {@code rawValue}, resolves the per-version compressor from {@code metadata} (including any ZSTD dictionary
802+
* cached on prior refreshes), decompresses the remainder, and runs the existing deserialization pipeline at the
803+
* embedded schema id. Reuses {@link #getDataRecordDeserializer} + {@link #tryToDeserialize} so any future change
804+
* to the in-band read path propagates here automatically.
805+
*/
806+
@Override
807+
public V decompressAndDeserialize(ByteBuffer rawValue, int version, K key) throws VeniceClientException {
808+
if (rawValue == null) {
809+
throw new IllegalArgumentException("rawValue must not be null");
810+
}
811+
if (rawValue.remaining() < Integer.BYTES) {
812+
throw new IllegalArgumentException(
813+
"rawValue must hold a 4-byte writer-schema-id prefix; got " + rawValue.remaining() + " bytes");
814+
}
815+
// Work on a duplicate so reading the schemaId prefix + decompressing does not advance the caller's buffer
816+
// position. duplicate() shares the underlying bytes but gives us our own position/limit/mark.
817+
ByteBuffer view = rawValue.duplicate();
818+
int schemaId = view.getInt();
819+
CompressionStrategy strategy = metadata.getCompressionStrategy(version);
820+
VeniceCompressor compressor = metadata.getCompressor(strategy, version);
821+
ByteBuffer decompressed;
822+
try {
823+
decompressed = compressor.decompress(view);
824+
} catch (IOException e) {
825+
throw new VeniceClientException("Failed to decompress value bytes for store: " + getStoreName(), e);
826+
}
827+
return tryToDeserialize(getDataRecordDeserializer(schemaId), decompressed, schemaId, key);
828+
}
829+
830+
@Override
831+
public void registerVersionSwitchListener(StoreVersionSwitchListener listener) {
832+
metadata.registerVersionSwitchListener(listener);
833+
}
834+
835+
@Override
836+
public void registerStoreConfigChangeListener(StoreConfigChangeListener listener) {
837+
metadata.registerStoreConfigChangeListener(listener);
838+
}
795839
}

clients/venice-client/src/main/java/com/linkedin/venice/fastclient/factory/ClientFactory.java

Lines changed: 27 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,21 @@ public static <K, V> AvroGenericStoreClient<K, V> getAndStartGenericStoreClient(
3232
return getAndStartGenericStoreClient(storeMetadata, clientConfig);
3333
}
3434

35+
/**
36+
* Builds the Fast Client wrapper chain without calling {@link AvroGenericStoreClient#start()}. Callers that need
37+
* to register a {@code StoreVersionSwitchListener} or {@code StoreConfigChangeListener} before the first metadata
38+
* refresh fires should use this entry point, register on the returned client, then call {@code start()}.
39+
*
40+
* <p>{@code register*Listener} only delivers transitions observed after the listener exists in the registry — a
41+
* listener registered after {@code start()} returns will not see the initial {@code (-1 -> currentVersion)}
42+
* transition committed by the first refresh. This factory plus {@code start()} from the caller is the supported
43+
* way to observe that initial transition.
44+
*/
45+
public static <K, V> AvroGenericStoreClient<K, V> getGenericStoreClient(ClientConfig clientConfig) {
46+
StoreMetadata storeMetadata = constructStoreMetadataReader(clientConfig);
47+
return buildGenericStoreClient(storeMetadata, clientConfig);
48+
}
49+
3550
public static <K, V extends SpecificRecord> AvroSpecificStoreClient<K, V> getAndStartSpecificStoreClient(
3651
ClientConfig clientConfig) {
3752
/**
@@ -63,6 +78,18 @@ private static StoreMetadata constructStoreMetadataReader(ClientConfig clientCon
6378
public static <K, V> AvroGenericStoreClient<K, V> getAndStartGenericStoreClient(
6479
StoreMetadata storeMetadata,
6580
ClientConfig clientConfig) {
81+
AvroGenericStoreClient<K, V> client = buildGenericStoreClient(storeMetadata, clientConfig);
82+
client.start();
83+
return client;
84+
}
85+
86+
/**
87+
* Builds the generic Fast Client wrapper chain over an externally-supplied {@link StoreMetadata} without calling
88+
* {@link AvroGenericStoreClient#start()}. See {@link #getGenericStoreClient(ClientConfig)} for the rationale.
89+
*/
90+
private static <K, V> AvroGenericStoreClient<K, V> buildGenericStoreClient(
91+
StoreMetadata storeMetadata,
92+
ClientConfig clientConfig) {
6693
final DispatchingAvroGenericStoreClient<K, V> dispatchingStoreClient = clientConfig.isVsonStore()
6794
? new DispatchingVsonStoreClient<>(storeMetadata, clientConfig)
6895
: new DispatchingAvroGenericStoreClient<>(storeMetadata, clientConfig);
@@ -92,8 +119,6 @@ public static <K, V> AvroGenericStoreClient<K, V> getAndStartGenericStoreClient(
92119
if (clientConfig.isDualReadEnabled()) {
93120
dualReadClient = new DualReadAvroGenericStoreClient<>(statsStoreClient, clientConfig);
94121
}
95-
dualReadClient.start();
96-
97122
return dualReadClient;
98123
}
99124

clients/venice-client/src/main/java/com/linkedin/venice/fastclient/meta/AbstractStoreMetadata.java

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
package com.linkedin.venice.fastclient.meta;
22

33
import com.linkedin.venice.client.exceptions.VeniceClientException;
4+
import com.linkedin.venice.client.store.listeners.StoreConfigChangeListener;
5+
import com.linkedin.venice.client.store.listeners.StoreConfigSnapshot;
6+
import com.linkedin.venice.client.store.listeners.StoreVersionSwitchListener;
47
import com.linkedin.venice.client.store.transport.TransportClientResponse;
58
import com.linkedin.venice.compression.CompressionStrategy;
69
import com.linkedin.venice.compression.CompressorFactory;
@@ -24,6 +27,7 @@
2427
import java.util.Map;
2528
import java.util.Set;
2629
import java.util.concurrent.CompletableFuture;
30+
import java.util.concurrent.CopyOnWriteArrayList;
2731
import java.util.concurrent.atomic.AtomicLong;
2832
import org.apache.logging.log4j.LogManager;
2933
import org.apache.logging.log4j.Logger;
@@ -36,6 +40,9 @@ public abstract class AbstractStoreMetadata implements StoreMetadata {
3640
private final InstanceHealthMonitor instanceHealthMonitor;
3741
protected volatile AbstractClientRoutingStrategy routingStrategy;
3842
protected final String storeName;
43+
private final CopyOnWriteArrayList<StoreVersionSwitchListener> versionSwitchListeners = new CopyOnWriteArrayList<>();
44+
private final CopyOnWriteArrayList<StoreConfigChangeListener> storeConfigChangeListeners =
45+
new CopyOnWriteArrayList<>();
3946

4047
public AbstractStoreMetadata(ClientConfig clientConfig) {
4148
this.clientConfig = clientConfig;
@@ -250,4 +257,90 @@ public <K> void routeRequest(RequestContext requestContext, RecordSerializer<K>
250257
throw new VeniceClientException("Unknown request type: " + requestType);
251258
}
252259

260+
@Override
261+
public void registerVersionSwitchListener(StoreVersionSwitchListener listener) {
262+
if (listener == null) {
263+
throw new IllegalArgumentException("StoreVersionSwitchListener must not be null");
264+
}
265+
versionSwitchListeners.addIfAbsent(listener);
266+
}
267+
268+
/**
269+
* Return a {@link Runnable} that will notify the registered version-switch listeners of the supplied
270+
* {@code (previousVersion, newVersion)} transition.
271+
*
272+
* <p>Build the {@link Runnable} inside the synchronized region that observed the transition, and run it after
273+
* the region has been exited so that listeners may safely call back into this metadata.
274+
*/
275+
protected Runnable buildVersionSwitchCallback(int previousVersion, int newVersion) {
276+
return () -> fireVersionSwitch(previousVersion, newVersion);
277+
}
278+
279+
/**
280+
* Notify the registered version-switch listeners of a transition. Each listener is invoked in registration order;
281+
* exceptions are caught and logged so that one failing listener cannot prevent others from running.
282+
*/
283+
protected void fireVersionSwitch(int previousVersion, int newVersion) {
284+
for (StoreVersionSwitchListener listener: versionSwitchListeners) {
285+
try {
286+
listener.onVersionSwitch(previousVersion, newVersion);
287+
} catch (Throwable t) {
288+
LOGGER.error(
289+
"Store {} version-switch listener {} threw on transition {} -> {}",
290+
storeName,
291+
listener.getClass().getName(),
292+
previousVersion,
293+
newVersion,
294+
t);
295+
}
296+
}
297+
}
298+
299+
@Override
300+
public void registerStoreConfigChangeListener(StoreConfigChangeListener listener) {
301+
if (listener == null) {
302+
throw new IllegalArgumentException("StoreConfigChangeListener must not be null");
303+
}
304+
storeConfigChangeListeners.addIfAbsent(listener);
305+
}
306+
307+
/**
308+
* Return a {@link Runnable} that will notify the registered store-config-change listeners of the supplied
309+
* {@code (previous, current)} transition.
310+
*/
311+
protected Runnable buildStoreConfigChangeCallback(StoreConfigSnapshot previous, StoreConfigSnapshot current) {
312+
return () -> fireStoreConfigChange(previous, current);
313+
}
314+
315+
/**
316+
* Notify the registered store-config-change listeners of a transition, but only if {@code previous} and
317+
* {@code current} differ — so callers may invoke this on every refresh without diffing first. Each listener is
318+
* invoked in registration order; exceptions are caught and logged so that one failing listener cannot prevent
319+
* others from running.
320+
*
321+
* @param previous prior snapshot, or {@code null} if no prior snapshot existed
322+
* @param current new snapshot; must not be {@code null}
323+
* @throws IllegalArgumentException if {@code current} is {@code null}
324+
*/
325+
protected void fireStoreConfigChange(StoreConfigSnapshot previous, StoreConfigSnapshot current) {
326+
if (current == null) {
327+
throw new IllegalArgumentException("current snapshot must not be null");
328+
}
329+
if (current.equals(previous)) {
330+
return;
331+
}
332+
for (StoreConfigChangeListener listener: storeConfigChangeListeners) {
333+
try {
334+
listener.onStoreConfigChange(previous, current);
335+
} catch (Throwable t) {
336+
LOGGER.error(
337+
"Store {} store-config-change listener {} threw on transition {} -> {}",
338+
storeName,
339+
listener.getClass().getName(),
340+
previous,
341+
current,
342+
t);
343+
}
344+
}
345+
}
253346
}

0 commit comments

Comments
 (0)