Skip to content

Commit 413d87d

Browse files
authored
[router][thin-client] Add getAllStoreNames capability via /stores endpoint (linkedin#2704)
Add GET /stores router endpoint backed by HelixReadOnlyStoreConfigRepository to list all non-system store names cluster-agnostically. Expose via StoreMetadataFetcher interface + RouterBasedStoreMetadataFetcher using D2TransportClient directly, and ClientFactory.createStoreMetadataFetcher().
1 parent cdc65fd commit 413d87d

14 files changed

Lines changed: 475 additions & 5 deletions

File tree

clients/venice-thin-client/src/main/java/com/linkedin/venice/client/store/ClientFactory.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,10 @@ public static StoreSchemaFetcher createStoreSchemaFetcher(ClientConfig clientCon
152152
new AvroGenericStoreClientImpl<>(getTransportClient(clientConfig), false, clientConfig));
153153
}
154154

155+
public static StoreMetadataFetcher createStoreMetadataFetcher(ClientConfig clientConfig) {
156+
return new RouterBasedStoreMetadataFetcher(clientConfig.getD2Client(), clientConfig.getD2ServiceName());
157+
}
158+
155159
private static D2TransportClient generateD2TransportClient(ClientConfig clientConfig) {
156160
String d2ServiceName = clientConfig.getD2ServiceName();
157161

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,100 @@
1+
package com.linkedin.venice.client.store;
2+
3+
import com.fasterxml.jackson.databind.DeserializationFeature;
4+
import com.fasterxml.jackson.databind.ObjectMapper;
5+
import com.linkedin.d2.balancer.D2Client;
6+
import com.linkedin.venice.client.store.transport.D2TransportClient;
7+
import com.linkedin.venice.client.store.transport.TransportClient;
8+
import com.linkedin.venice.client.store.transport.TransportClientResponse;
9+
import com.linkedin.venice.controllerapi.MultiStoreResponse;
10+
import com.linkedin.venice.exceptions.VeniceException;
11+
import com.linkedin.venice.utils.ObjectMapperFactory;
12+
import java.io.IOException;
13+
import java.util.Arrays;
14+
import java.util.Collections;
15+
import java.util.HashSet;
16+
import java.util.Set;
17+
import java.util.concurrent.ExecutionException;
18+
import java.util.concurrent.TimeUnit;
19+
import java.util.concurrent.TimeoutException;
20+
21+
22+
/**
23+
* Router-based implementation for fetching store metadata that is not cluster-specific.
24+
* Unlike {@link com.linkedin.venice.client.schema.RouterBasedStoreSchemaFetcher}, this class
25+
* is not tied to a specific store and operates on metadata available globally across clusters.
26+
*
27+
* This class uses {@link D2TransportClient} directly (rather than {@link AbstractAvroStoreClient})
28+
* to avoid store-level D2 service discovery, which requires a store name and is unnecessary for
29+
* cluster-agnostic endpoints like {@code /stores}.
30+
*/
31+
public class RouterBasedStoreMetadataFetcher implements StoreMetadataFetcher {
32+
public static final String TYPE_STORES = "stores";
33+
34+
private static final ObjectMapper OBJECT_MAPPER = ObjectMapperFactory.getInstance();
35+
36+
// Ignore unknown fields while parsing json response.
37+
static {
38+
OBJECT_MAPPER.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
39+
}
40+
41+
static final long GET_TIMEOUT_IN_SECONDS = 30;
42+
43+
private final TransportClient transportClient;
44+
45+
public RouterBasedStoreMetadataFetcher(D2Client d2Client, String d2ServiceName) {
46+
this.transportClient = new D2TransportClient(d2ServiceName, d2Client);
47+
}
48+
49+
// VisibleForTesting
50+
RouterBasedStoreMetadataFetcher(TransportClient transportClient) {
51+
this.transportClient = transportClient;
52+
}
53+
54+
/**
55+
* Returns all store names available across all clusters, as seen by the router's
56+
* non-cluster-specific {@link com.linkedin.venice.helix.HelixReadOnlyStoreConfigRepository}.
57+
*/
58+
@Override
59+
public Set<String> getAllStoreNames() {
60+
byte[] responseBody;
61+
try {
62+
TransportClientResponse response =
63+
transportClient.get(TYPE_STORES, Collections.emptyMap()).get(GET_TIMEOUT_IN_SECONDS, TimeUnit.SECONDS);
64+
if (response == null) {
65+
throw new VeniceException("Received null response from router for path: " + TYPE_STORES);
66+
}
67+
responseBody = response.getBody();
68+
if (responseBody == null) {
69+
throw new VeniceException("Received empty response body from router for path: " + TYPE_STORES);
70+
}
71+
} catch (ExecutionException | InterruptedException | TimeoutException e) {
72+
if (e instanceof InterruptedException) {
73+
Thread.currentThread().interrupt();
74+
}
75+
throw new VeniceException("Failed to fetch store names from router for path: " + TYPE_STORES, e);
76+
}
77+
78+
MultiStoreResponse multiStoreResponse;
79+
try {
80+
multiStoreResponse = OBJECT_MAPPER.readValue(responseBody, MultiStoreResponse.class);
81+
} catch (IOException e) {
82+
throw new VeniceException("Failed to deserialize store names response", e);
83+
}
84+
85+
if (multiStoreResponse.isError()) {
86+
throw new VeniceException("Received error while fetching store names: " + multiStoreResponse.getError());
87+
}
88+
89+
String[] stores = multiStoreResponse.getStores();
90+
if (stores == null) {
91+
return Collections.emptySet();
92+
}
93+
return new HashSet<>(Arrays.asList(stores));
94+
}
95+
96+
@Override
97+
public void close() throws IOException {
98+
transportClient.close();
99+
}
100+
}
Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
package com.linkedin.venice.client.store;
2+
3+
import java.io.Closeable;
4+
import java.util.Set;
5+
6+
7+
/**
8+
* Public interface for fetching Venice store metadata that is not tied to a specific store or cluster.
9+
* It is intended to provide metadata operations that span across all stores and clusters, as opposed to
10+
* {@link com.linkedin.venice.client.schema.StoreSchemaFetcher} which is scoped to a single store.
11+
*/
12+
public interface StoreMetadataFetcher extends Closeable {
13+
/**
14+
* Returns all store names available across all clusters.
15+
*/
16+
Set<String> getAllStoreNames();
17+
}
Lines changed: 130 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,130 @@
1+
package com.linkedin.venice.client.store;
2+
3+
import static com.linkedin.venice.client.store.RouterBasedStoreMetadataFetcher.TYPE_STORES;
4+
5+
import com.fasterxml.jackson.databind.ObjectMapper;
6+
import com.linkedin.venice.client.store.transport.TransportClient;
7+
import com.linkedin.venice.client.store.transport.TransportClientResponse;
8+
import com.linkedin.venice.controllerapi.MultiStoreResponse;
9+
import com.linkedin.venice.exceptions.VeniceException;
10+
import com.linkedin.venice.utils.ObjectMapperFactory;
11+
import java.io.IOException;
12+
import java.util.Set;
13+
import java.util.concurrent.CompletableFuture;
14+
import java.util.concurrent.ExecutionException;
15+
import java.util.concurrent.TimeoutException;
16+
import org.mockito.Mockito;
17+
import org.testng.Assert;
18+
import org.testng.annotations.Test;
19+
20+
21+
public class RouterBasedStoreMetadataFetcherTest {
22+
private static final ObjectMapper OBJECT_MAPPER = ObjectMapperFactory.getInstance();
23+
24+
private TransportClient mockTransportClientWith(MultiStoreResponse response) throws Exception {
25+
TransportClient mockClient = Mockito.mock(TransportClient.class);
26+
byte[] responseBytes = OBJECT_MAPPER.writeValueAsBytes(response);
27+
TransportClientResponse transportResponse = Mockito.mock(TransportClientResponse.class);
28+
Mockito.doReturn(responseBytes).when(transportResponse).getBody();
29+
CompletableFuture<TransportClientResponse> mockFuture = CompletableFuture.completedFuture(transportResponse);
30+
Mockito.doReturn(mockFuture).when(mockClient).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
31+
return mockClient;
32+
}
33+
34+
@Test
35+
public void testGetAllStoreNames() throws Exception {
36+
MultiStoreResponse multiStoreResponse = new MultiStoreResponse();
37+
multiStoreResponse.setStores(new String[] { "store_a", "store_b", "store_c" });
38+
39+
TransportClient mockClient = mockTransportClientWith(multiStoreResponse);
40+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
41+
42+
Set<String> storeNames = fetcher.getAllStoreNames();
43+
44+
Assert.assertEquals(storeNames.size(), 3);
45+
Assert.assertTrue(storeNames.contains("store_a"));
46+
Assert.assertTrue(storeNames.contains("store_b"));
47+
Assert.assertTrue(storeNames.contains("store_c"));
48+
Mockito.verify(mockClient, Mockito.times(1)).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
49+
}
50+
51+
@Test
52+
public void testGetAllStoreNamesReturnsEmptyWhenNullStores() throws Exception {
53+
MultiStoreResponse multiStoreResponse = new MultiStoreResponse();
54+
multiStoreResponse.setStores(null);
55+
56+
TransportClient mockClient = mockTransportClientWith(multiStoreResponse);
57+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
58+
59+
Set<String> storeNames = fetcher.getAllStoreNames();
60+
61+
Assert.assertTrue(storeNames.isEmpty());
62+
}
63+
64+
@Test
65+
public void testGetAllStoreNamesThrowsOnRouterError() throws Exception {
66+
MultiStoreResponse errorResponse = new MultiStoreResponse();
67+
errorResponse.setError("Internal server error");
68+
69+
TransportClient mockClient = mockTransportClientWith(errorResponse);
70+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
71+
72+
Assert.assertThrows(VeniceException.class, fetcher::getAllStoreNames);
73+
}
74+
75+
@Test
76+
public void testGetAllStoreNamesThrowsOnNullResponse() throws Exception {
77+
TransportClient mockClient = Mockito.mock(TransportClient.class);
78+
CompletableFuture<TransportClientResponse> mockFuture = CompletableFuture.completedFuture(null);
79+
Mockito.doReturn(mockFuture).when(mockClient).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
80+
81+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
82+
83+
Assert.assertThrows(VeniceException.class, fetcher::getAllStoreNames);
84+
}
85+
86+
@Test
87+
public void testGetAllStoreNamesThrowsOnNullResponseBody() throws Exception {
88+
TransportClient mockClient = Mockito.mock(TransportClient.class);
89+
TransportClientResponse transportResponse = Mockito.mock(TransportClientResponse.class);
90+
Mockito.doReturn(null).when(transportResponse).getBody();
91+
CompletableFuture<TransportClientResponse> mockFuture = CompletableFuture.completedFuture(transportResponse);
92+
Mockito.doReturn(mockFuture).when(mockClient).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
93+
94+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
95+
96+
Assert.assertThrows(VeniceException.class, fetcher::getAllStoreNames);
97+
}
98+
99+
@Test
100+
public void testGetAllStoreNamesThrowsOnTransportFailure() throws Exception {
101+
TransportClient mockClient = Mockito.mock(TransportClient.class);
102+
CompletableFuture<TransportClientResponse> failedFuture = new CompletableFuture<>();
103+
failedFuture.completeExceptionally(new ExecutionException("transport error", new RuntimeException()));
104+
Mockito.doReturn(failedFuture).when(mockClient).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
105+
106+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
107+
108+
Assert.assertThrows(VeniceException.class, fetcher::getAllStoreNames);
109+
}
110+
111+
@Test
112+
public void testGetAllStoreNamesThrowsOnTimeout() throws Exception {
113+
TransportClient mockClient = Mockito.mock(TransportClient.class);
114+
CompletableFuture<TransportClientResponse> stalledFuture = new CompletableFuture<>();
115+
stalledFuture.completeExceptionally(new TimeoutException("timed out"));
116+
Mockito.doReturn(stalledFuture).when(mockClient).get(Mockito.eq(TYPE_STORES), Mockito.anyMap());
117+
118+
StoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
119+
120+
Assert.assertThrows(VeniceException.class, fetcher::getAllStoreNames);
121+
}
122+
123+
@Test
124+
public void testClose() throws IOException {
125+
TransportClient mockClient = Mockito.mock(TransportClient.class);
126+
RouterBasedStoreMetadataFetcher fetcher = new RouterBasedStoreMetadataFetcher(mockClient);
127+
fetcher.close();
128+
Mockito.verify(mockClient, Mockito.times(1)).close();
129+
}
130+
}

internal/venice-common/src/main/java/com/linkedin/venice/controllerapi/MultiStoreResponse.java renamed to internal/venice-client-common/src/main/java/com/linkedin/venice/controllerapi/MultiStoreResponse.java

File renamed without changes.

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

Lines changed: 30 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import com.linkedin.venice.exceptions.VeniceException;
77
import com.linkedin.venice.exceptions.VeniceNoStoreException;
88
import com.linkedin.venice.meta.ReadOnlyStoreConfigRepository;
9+
import com.linkedin.venice.meta.Store;
910
import com.linkedin.venice.meta.StoreConfig;
1011
import com.linkedin.venice.utils.concurrent.VeniceConcurrentHashMap;
1112
import java.util.Collections;
@@ -34,6 +35,7 @@ public class HelixReadOnlyStoreConfigRepository implements ReadOnlyStoreConfigRe
3435

3536
private final Map<String, StoreConfig> loadedStoreConfigMap;
3637
private final AtomicReference<Set<String>> availableStoreSet;
38+
private final AtomicReference<Set<String>> availableRegularStoreSet;
3739
private final ZkStoreConfigAccessor accessor;
3840
private final StoreConfigChangedListener storeConfigChangedListener;
3941
private final StoreConfigAddedOrDeletedChangedListener storeConfigAddedOrDeletedListener;
@@ -49,6 +51,7 @@ public HelixReadOnlyStoreConfigRepository(ZkClient zkClient, ZkStoreConfigAccess
4951
this.accessor = accessor;
5052
this.loadedStoreConfigMap = new VeniceConcurrentHashMap<>();
5153
this.availableStoreSet = new AtomicReference<>(new HashSet<>());
54+
this.availableRegularStoreSet = new AtomicReference<>(new HashSet<>());
5255
storeConfigChangedListener = new StoreConfigChangedListener();
5356
storeConfigAddedOrDeletedListener = new StoreConfigAddedOrDeletedChangedListener();
5457
// This repository already retry on getChildren, so do not need extra retry in listener.
@@ -62,8 +65,10 @@ public HelixReadOnlyStoreConfigRepository(ZkClient zkClient, ZkStoreConfigAccess
6265
public void refresh() {
6366
LOGGER.info("Loading all store names from zk.");
6467
accessor.subscribeStoreConfigAddedOrDeletedListener(storeConfigAddedOrDeletedListener);
65-
availableStoreSet.set(new HashSet<>(accessor.getAllStores()));
66-
LOGGER.info("Found {} stores.", availableStoreSet.get().size());
68+
Set<String> allStores = new HashSet<>(accessor.getAllStores());
69+
availableStoreSet.set(allStores);
70+
availableRegularStoreSet.set(filterRegularStores(allStores));
71+
LOGGER.info("Found {} stores.", allStores.size());
6772
zkClient.subscribeStateChanges(zkStateListener);
6873
LOGGER.info("All store names are loaded.");
6974
}
@@ -77,6 +82,7 @@ public void clear() {
7782
}
7883
this.loadedStoreConfigMap.clear();
7984
this.availableStoreSet.set(Collections.emptySet());
85+
this.availableRegularStoreSet.set(Collections.emptySet());
8086
zkClient.unsubscribeStateChanges(zkStateListener);
8187
LOGGER.info("Cleared all store configs in local");
8288
}
@@ -131,6 +137,24 @@ public StoreConfig getStoreConfigOrThrow(String storeName) {
131137
return storeConfig.get();
132138
}
133139

140+
@Override
141+
public Set<String> getStores(boolean includeSystemStores) {
142+
if (includeSystemStores) {
143+
return Collections.unmodifiableSet(getAvailableStoreSet());
144+
}
145+
return Collections.unmodifiableSet(availableRegularStoreSet.get());
146+
}
147+
148+
private static Set<String> filterRegularStores(Set<String> stores) {
149+
Set<String> regularStores = new HashSet<>();
150+
for (String storeName: stores) {
151+
if (!storeName.startsWith(Store.SYSTEM_STORE_NAME_PREFIX)) {
152+
regularStores.add(storeName);
153+
}
154+
}
155+
return regularStores;
156+
}
157+
134158
@VisibleForTesting
135159
StoreConfigAddedOrDeletedChangedListener getStoreConfigAddedOrDeletedListener() {
136160
return storeConfigAddedOrDeletedListener;
@@ -165,8 +189,10 @@ public void handleChildChange(String parentPath, List<String> currentChildren) {
165189
newStoresCount,
166190
storeSetSnapshot.size());
167191

168-
// update the available store set
169-
availableStoreSet.set(new HashSet<>(currentChildren));
192+
// update the available store set and the cached regular store set
193+
Set<String> newStoreSet = new HashSet<>(currentChildren);
194+
availableStoreSet.set(newStoreSet);
195+
availableRegularStoreSet.set(filterRegularStores(newStoreSet));
170196

171197
// Deleted store configs
172198
for (String deletedStore: storeSetSnapshot) {

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import com.linkedin.venice.meta.StoreConfig;
66
import com.linkedin.venice.views.VeniceView;
77
import java.util.Optional;
8+
import java.util.Set;
89

910

1011
/**
@@ -37,4 +38,9 @@ public Optional<StoreConfig> getStoreConfig(String storeName) {
3738
public StoreConfig getStoreConfigOrThrow(String storeName) {
3839
return storeConfigRepository.getStoreConfigOrThrow(VeniceView.getStoreName(storeName));
3940
}
41+
42+
@Override
43+
public Set<String> getStores(boolean includeSystemStores) {
44+
return storeConfigRepository.getStores(includeSystemStores);
45+
}
4046
}

internal/venice-common/src/main/java/com/linkedin/venice/meta/ReadOnlyStoreConfigRepository.java

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

33
import java.util.Optional;
4+
import java.util.Set;
45

56

67
/**
@@ -10,4 +11,6 @@ public interface ReadOnlyStoreConfigRepository {
1011
Optional<StoreConfig> getStoreConfig(String storeName);
1112

1213
StoreConfig getStoreConfigOrThrow(String storeName);
14+
15+
Set<String> getStores(boolean includeSystemStores);
1316
}

0 commit comments

Comments
 (0)