Skip to content

Commit c7c55d9

Browse files
authored
Remove offset fields from AdminMetadata V2 format (linkedin#2465)
Remove dual support for V1 (numeric offset) and V2 (PubSubPosition) admin topic metadata and standardize fully on V2 AdminMetadata format. Drop USE_V2_ADMIN_TOPIC_METADATA config, V1 ZK path, cluster config toggle, V1 methods in ZkAdminTopicMetadataAccessor and AdminMetadata, and failed_admin_message_offset metric.
1 parent c150276 commit c7c55d9

26 files changed

Lines changed: 552 additions & 694 deletions

clients/venice-admin-tool/src/test/java/com/linkedin/venice/TestTreeNode.java

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -66,18 +66,18 @@ public void testPathsTreeToList() {
6666
Assert.assertTrue(list.contains("/venice-parent/storeConfigs"));
6767
Assert.assertTrue(list.contains("/venice-parent/cluster1"));
6868
Assert.assertFalse(list.contains("/venice-parent/cluster1/storeConfigs"));
69-
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadata"));
70-
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadata/file1"));
71-
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadata/file2"));
72-
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadata/file2/file3"));
69+
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadataV2"));
70+
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadataV2/file1"));
71+
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadataV2/file2"));
72+
Assert.assertTrue(list.contains("/venice-parent/cluster1/adminTopicMetadataV2/file2/file3"));
7373
Assert.assertTrue(list.contains("/venice-parent/cluster1/executionids"));
7474
Assert.assertTrue(list.contains("/venice-parent/cluster1/ParentOfflinePushes"));
7575
Assert.assertTrue(list.contains("/venice-parent/cluster1/routers"));
7676
Assert.assertTrue(list.contains("/venice-parent/cluster1/StoreGraveyard"));
7777
Assert.assertTrue(list.contains("/venice-parent/cluster1/Stores"));
7878
Assert.assertTrue(list.contains("/venice-parent/cluster2"));
7979
Assert.assertFalse(list.contains("/venice-parent/cluster2/storeConfigs"));
80-
Assert.assertTrue(list.contains("/venice-parent/cluster2/adminTopicMetadata"));
80+
Assert.assertTrue(list.contains("/venice-parent/cluster2/adminTopicMetadataV2"));
8181
Assert.assertTrue(list.contains("/venice-parent/cluster2/executionids"));
8282
Assert.assertTrue(list.contains("/venice-parent/cluster2/ParentOfflinePushes"));
8383
Assert.assertTrue(list.contains("/venice-parent/cluster2/routers"));
@@ -156,16 +156,16 @@ private List<String> getPaths() {
156156
List<String> paths = new ArrayList<>();
157157
paths.add("/venice-parent/storeConfigs");
158158
paths.add("/venice-parent/cluster1");
159-
paths.add("/venice-parent/cluster1/adminTopicMetadata");
160-
paths.add("/venice-parent/cluster1/adminTopicMetadata/file1");
161-
paths.add("/venice-parent/cluster1/adminTopicMetadata/file2/file3"); // tests both /file2 and /file2/file3
159+
paths.add("/venice-parent/cluster1/adminTopicMetadataV2");
160+
paths.add("/venice-parent/cluster1/adminTopicMetadataV2/file1");
161+
paths.add("/venice-parent/cluster1/adminTopicMetadataV2/file2/file3"); // tests both /file2 and /file2/file3
162162
paths.add("/venice-parent/cluster1/executionids");
163163
paths.add("/venice-parent/cluster1/ParentOfflinePushes");
164164
paths.add("/venice-parent/cluster1/routers");
165165
paths.add("/venice-parent/cluster1/StoreGraveyard");
166166
paths.add("/venice-parent/cluster1/Stores");
167167
paths.add("/venice-parent/cluster2");
168-
paths.add("/venice-parent/cluster2/adminTopicMetadata");
168+
paths.add("/venice-parent/cluster2/adminTopicMetadataV2");
169169
paths.add("/venice-parent/cluster2/executionids");
170170
paths.add("/venice-parent/cluster2/ParentOfflinePushes");
171171
paths.add("/venice-parent/cluster2/routers");

clients/venice-admin-tool/src/test/java/com/linkedin/venice/TestZkCopier.java

Lines changed: 21 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -62,22 +62,22 @@ public void testGetVenicePathsFromList() {
6262
Assert.assertFalse(venicePaths.contains("/venice/storeConfigs"));
6363
Assert.assertFalse(venicePaths.contains("/venice/cluster1"));
6464
Assert.assertFalse(venicePaths.contains("/venice/cluster1/storeConfigs"));
65-
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadata"));
66-
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadata/file1"));
67-
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadata/file2/file3"));
65+
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadataV2"));
66+
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadataV2/file1"));
67+
Assert.assertFalse(venicePaths.contains("/venice/cluster1/adminTopicMetadataV2/file2/file3"));
6868
Assert.assertFalse(venicePaths.contains("/venice-parent/cluster1/storeConfigs"));
69-
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadata/file1"));
70-
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadata/file2"));
71-
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadata/file2/file3"));
69+
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadataV2/file1"));
70+
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadataV2/file2"));
71+
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadataV2/file2/file3"));
7272
Assert.assertFalse(venicePaths.contains("/venice-parent/cluster2/storeConfigs"));
7373
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/executionids/file1"));
7474
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/executionids/file2"));
7575
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/executionids/file2/file3"));
7676
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster"));
7777
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/storeConfigs"));
78-
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadata"));
79-
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadata/file1"));
80-
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadata/file2/file3"));
78+
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadataV2"));
79+
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadataV2/file1"));
80+
Assert.assertFalse(venicePaths.contains("/venice-parent/helix-cluster/adminTopicMetadataV2/file2/file3"));
8181
testVenicePathsContainsAsserts(venicePaths);
8282
}
8383

@@ -116,14 +116,14 @@ private void testVenicePathsContainsAsserts(List<String> venicePaths) {
116116
Assert.assertTrue(venicePaths.contains("/venice-parent"));
117117
Assert.assertTrue(venicePaths.contains("/venice-parent/storeConfigs"));
118118
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1"));
119-
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadata"));
119+
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/adminTopicMetadataV2"));
120120
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/executionids"));
121121
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/ParentOfflinePushes"));
122122
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/routers"));
123123
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/StoreGraveyard"));
124124
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster1/Stores"));
125125
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2"));
126-
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/adminTopicMetadata"));
126+
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/adminTopicMetadataV2"));
127127
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/executionids"));
128128
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/ParentOfflinePushes"));
129129
Assert.assertTrue(venicePaths.contains("/venice-parent/cluster2/routers"));
@@ -146,24 +146,24 @@ private List<String> getPaths() {
146146
zkPaths.add("/venice/storeConfigs");
147147
zkPaths.add("/venice/cluster1");
148148
zkPaths.add("/venice/cluster1/storeConfigs");
149-
zkPaths.add("/venice/cluster1/adminTopicMetadata");
150-
zkPaths.add("/venice/cluster1/adminTopicMetadata/file1");
151-
zkPaths.add("/venice/cluster1/adminTopicMetadata/file2/file3");
149+
zkPaths.add("/venice/cluster1/adminTopicMetadataV2");
150+
zkPaths.add("/venice/cluster1/adminTopicMetadataV2/file1");
151+
zkPaths.add("/venice/cluster1/adminTopicMetadataV2/file2/file3");
152152
zkPaths.add("/venice-parent");
153153
zkPaths.add("/venice-parent/storeConfigs");
154154
zkPaths.add("/venice-parent/cluster1");
155155
zkPaths.add("/venice-parent/cluster1/storeConfigs");
156-
zkPaths.add("/venice-parent/cluster1/adminTopicMetadata");
157-
zkPaths.add("/venice-parent/cluster1/adminTopicMetadata/file1");
158-
zkPaths.add("/venice-parent/cluster1/adminTopicMetadata/file2/file3");
156+
zkPaths.add("/venice-parent/cluster1/adminTopicMetadataV2");
157+
zkPaths.add("/venice-parent/cluster1/adminTopicMetadataV2/file1");
158+
zkPaths.add("/venice-parent/cluster1/adminTopicMetadataV2/file2/file3");
159159
zkPaths.add("/venice-parent/cluster1/executionids");
160160
zkPaths.add("/venice-parent/cluster1/ParentOfflinePushes");
161161
zkPaths.add("/venice-parent/cluster1/routers");
162162
zkPaths.add("/venice-parent/cluster1/StoreGraveyard");
163163
zkPaths.add("/venice-parent/cluster1/Stores");
164164
zkPaths.add("/venice-parent/cluster2");
165165
zkPaths.add("/venice-parent/cluster2/storeConfigs");
166-
zkPaths.add("/venice-parent/cluster2/adminTopicMetadata");
166+
zkPaths.add("/venice-parent/cluster2/adminTopicMetadataV2");
167167
zkPaths.add("/venice-parent/cluster2/executionids");
168168
zkPaths.add("/venice-parent/cluster2/executionids/file1");
169169
zkPaths.add("/venice-parent/cluster2/executionids/file2/file3");
@@ -173,9 +173,9 @@ private List<String> getPaths() {
173173
zkPaths.add("/venice-parent/cluster2/Stores");
174174
zkPaths.add("/venice-parent/helix-cluster");
175175
zkPaths.add("/venice-parent/helix-cluster/storeConfigs");
176-
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadata");
177-
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadata/file1");
178-
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadata/file2/file3");
176+
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadataV2");
177+
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadataV2/file1");
178+
zkPaths.add("/venice-parent/helix-cluster/adminTopicMetadataV2/file2/file3");
179179
return zkPaths;
180180
}
181181
}

clients/venice-admin-tool/src/test/resources/venice_paths.txt

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,17 +3,17 @@
33
/venice-parent/cluster1/ParentOfflinePushes
44
/venice-parent/cluster1/StoreGraveyard
55
/venice-parent/cluster1/Stores
6-
/venice-parent/cluster1/adminTopicMetadata
7-
/venice-parent/cluster1/adminTopicMetadata/file1
8-
/venice-parent/cluster1/adminTopicMetadata/file2
9-
/venice-parent/cluster1/adminTopicMetadata/file2/file3
6+
/venice-parent/cluster1/adminTopicMetadataV2
7+
/venice-parent/cluster1/adminTopicMetadataV2/file1
8+
/venice-parent/cluster1/adminTopicMetadataV2/file2
9+
/venice-parent/cluster1/adminTopicMetadataV2/file2/file3
1010
/venice-parent/cluster1/executionids
1111
/venice-parent/cluster1/routers
1212
/venice-parent/cluster2
1313
/venice-parent/cluster2/ParentOfflinePushes
1414
/venice-parent/cluster2/StoreGraveyard
1515
/venice-parent/cluster2/Stores
16-
/venice-parent/cluster2/adminTopicMetadata
16+
/venice-parent/cluster2/adminTopicMetadataV2
1717
/venice-parent/cluster2/executionids
1818
/venice-parent/cluster2/executionids/file1
1919
/venice-parent/cluster2/executionids/file2

clients/venice-admin-tool/src/test/resources/zk_paths.txt

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -2,24 +2,24 @@
22
/venice/storeConfigs
33
/venice/cluster1
44
/venice/cluster1/storeConfigs
5-
/venice/cluster1/adminTopicMetadata
6-
/venice/cluster1/adminTopicMetadata/file1
7-
/venice/cluster1/adminTopicMetadata/file2/file3
5+
/venice/cluster1/adminTopicMetadataV2
6+
/venice/cluster1/adminTopicMetadataV2/file1
7+
/venice/cluster1/adminTopicMetadataV2/file2/file3
88
/venice-parent
99
/venice-parent/storeConfigs
1010
/venice-parent/cluster1
1111
/venice-parent/cluster1/storeConfigs
12-
/venice-parent/cluster1/adminTopicMetadata
13-
/venice-parent/cluster1/adminTopicMetadata/file1
14-
/venice-parent/cluster1/adminTopicMetadata/file2/file3
12+
/venice-parent/cluster1/adminTopicMetadataV2
13+
/venice-parent/cluster1/adminTopicMetadataV2/file1
14+
/venice-parent/cluster1/adminTopicMetadataV2/file2/file3
1515
/venice-parent/cluster1/executionids
1616
/venice-parent/cluster1/ParentOfflinePushes
1717
/venice-parent/cluster1/routers
1818
/venice-parent/cluster1/StoreGraveyard
1919
/venice-parent/cluster1/Stores
2020
/venice-parent/cluster2
2121
/venice-parent/cluster2/storeConfigs
22-
/venice-parent/cluster2/adminTopicMetadata
22+
/venice-parent/cluster2/adminTopicMetadataV2
2323
/venice-parent/cluster2/executionids
2424
/venice-parent/cluster2/executionids/file1
2525
/venice-parent/cluster2/executionids/file2/file3
@@ -29,6 +29,6 @@
2929
/venice-parent/cluster2/Stores
3030
/venice-parent/helix-cluster
3131
/venice-parent/helix-cluster/storeConfigs
32-
/venice-parent/helix-cluster/adminTopicMetadata
33-
/venice-parent/helix-cluster/adminTopicMetadata/file1
34-
/venice-parent/helix-cluster/adminTopicMetadata/file2/file3
32+
/venice-parent/helix-cluster/adminTopicMetadataV2
33+
/venice-parent/helix-cluster/adminTopicMetadataV2/file1
34+
/venice-parent/helix-cluster/adminTopicMetadataV2/file2/file3

internal/venice-common/src/main/java/com/linkedin/venice/ConfigKeys.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2974,7 +2974,6 @@ private ConfigKeys() {
29742974
"controller.enable.realtime.topic.versioning";
29752975

29762976
public static final boolean DEFAULT_CONTROLLER_ENABLE_REAL_TIME_TOPIC_VERSIONING = false;
2977-
public final static String USE_V2_ADMIN_TOPIC_METADATA = "controller.use.v2.admin.topic.metadata";
29782977
public static final String CONTROLLER_ENABLE_HYBRID_STORE_PARTITION_COUNT_UPDATE =
29792978
"controller.enable.hybrid.store.partition.count.update";
29802979
public static final String PUSH_JOB_VIEW_CONFIGS = "push.job.view.configs";

internal/venice-common/src/main/java/com/linkedin/venice/zk/VeniceZkPaths.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,7 @@
1010
* This class contains constants that represent Venice-managed ZooKeeper paths.
1111
*/
1212
public class VeniceZkPaths {
13-
public static final String ADMIN_TOPIC_METADATA = "adminTopicMetadata";
14-
// new admin topic metadata structure is incompatible with the old one, so creating a new "v2" path
15-
public static final String ADMIN_TOPIC_METADATA_V2 = "adminTopicMetadataV2";
13+
public static final String ADMIN_TOPIC_METADATA = "adminTopicMetadataV2";
1614
public static final String CLUSTER_CONFIG = "ClusterConfig";
1715
public static final String DARK_CLUSTER_CONFIG = "DarkClusterConfig";
1816
public static final String EXECUTION_IDS = "executionids";

internal/venice-test-common/src/main/java/com/linkedin/venice/admin/InMemoryAdminTopicMetadataAccessor.java

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -22,12 +22,6 @@ public void updateMetadata(String clusterName, AdminMetadata metadataDelta) {
2222
if (metadataDelta.getExecutionId() != null) {
2323
inMemoryMetadata.setExecutionId(metadataDelta.getExecutionId());
2424
}
25-
if (metadataDelta.getOffset() != null) {
26-
inMemoryMetadata.setOffset(metadataDelta.getOffset());
27-
}
28-
if (metadataDelta.getUpstreamOffset() != null) {
29-
inMemoryMetadata.setUpstreamOffset(metadataDelta.getUpstreamOffset());
30-
}
3125
if (!metadataDelta.getAdminOperationProtocolVersion().equals(UNDEFINED_VALUE)) {
3226
inMemoryMetadata.setAdminOperationProtocolVersion(metadataDelta.getAdminOperationProtocolVersion());
3327
}

internal/venice-test-common/src/main/java/com/linkedin/venice/pubsub/mock/InMemoryPubSubPositionFactory.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import static com.linkedin.venice.pubsub.mock.InMemoryPubSubPosition.*;
44

5+
import com.linkedin.venice.pubsub.PubSubPositionDeserializer;
56
import com.linkedin.venice.pubsub.PubSubPositionFactory;
67
import com.linkedin.venice.pubsub.PubSubPositionTypeRegistry;
78
import com.linkedin.venice.pubsub.api.PubSubPosition;
@@ -35,4 +36,8 @@ public static PubSubPositionTypeRegistry getPositionTypeRegistryWithInMemoryPosi
3536
typeIdToFactory.put(INMEMORY_PUBSUB_POSITION_TYPE_ID, InMemoryPubSubPositionFactory.class.getName());
3637
return new PubSubPositionTypeRegistry(typeIdToFactory);
3738
}
39+
40+
public static PubSubPositionDeserializer getPubSubPositionDeserializerWithInMemoryPosition() {
41+
return new PubSubPositionDeserializer(getPositionTypeRegistryWithInMemoryPosition());
42+
}
3843
}

services/venice-controller/src/main/java/com/linkedin/venice/controller/AdminTopicMetadataAccessor.java

Lines changed: 0 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -3,83 +3,34 @@
33
import com.linkedin.venice.controller.kafka.consumer.AdminMetadata;
44
import com.linkedin.venice.pubsub.api.PubSubPosition;
55
import com.linkedin.venice.utils.Pair;
6-
import java.util.HashMap;
7-
import java.util.Map;
8-
import java.util.Optional;
96

107

118
/**
129
* This class provides a set of methods to access and update metadata for admin topics.
1310
*/
1411
public abstract class AdminTopicMetadataAccessor {
15-
public static final String OFFSET_KEY = "offset";
1612
public static final String POSITION_KEY = "position";
1713
/**
1814
* When remote consumption is enabled, child controller will consume directly from the source admin topic; an extra
1915
* metadata called upstream offset will be maintained, which indicate the last offset in the source admin topic that
2016
* gets processed successfully.
2117
*/
22-
public static final String UPSTREAM_OFFSET_KEY = "upstreamOffset";
2318
public static final String UPSTREAM_POSITION_KEY = "upstreamPosition";
2419
public static final String EXECUTION_ID_KEY = "executionId";
2520
public static final String ADMIN_OPERATION_PROTOCOL_VERSION_KEY = "adminOperationProtocolVersion";
2621
public static final Long UNDEFINED_VALUE = -1L;
2722

28-
/**
29-
* @return a map with {@linkplain AdminTopicMetadataAccessor#OFFSET_KEY}, {@linkplain AdminTopicMetadataAccessor#UPSTREAM_OFFSET_KEY},
30-
* {@linkplain AdminTopicMetadataAccessor#EXECUTION_ID_KEY}, {@linkplain AdminTopicMetadataAccessor#ADMIN_OPERATION_PROTOCOL_VERSION_KEY}
31-
* specified to input values.
32-
*/
33-
public static Map<String, Long> generateMetadataMap(
34-
Optional<Long> localOffset,
35-
Optional<Long> upstreamOffset,
36-
Optional<Long> executionId,
37-
Optional<Long> adminOperationProtocolVersion) {
38-
Map<String, Long> metadata = new HashMap<>();
39-
localOffset.ifPresent(offset -> metadata.put(OFFSET_KEY, offset));
40-
upstreamOffset.ifPresent(offset -> metadata.put(UPSTREAM_OFFSET_KEY, offset));
41-
executionId.ifPresent(id -> metadata.put(EXECUTION_ID_KEY, id));
42-
adminOperationProtocolVersion.ifPresent(version -> metadata.put(ADMIN_OPERATION_PROTOCOL_VERSION_KEY, version));
43-
return metadata;
44-
}
45-
46-
/**
47-
* @return a pair of values to which the specified keys are mapped to {@linkplain AdminTopicMetadataAccessor#OFFSET_KEY}
48-
* and {@linkplain AdminTopicMetadataAccessor#UPSTREAM_OFFSET_KEY}.
49-
*/
50-
public static Pair<Long, Long> getOffsets(Map<String, Long> metadata) {
51-
long localOffset = metadata.getOrDefault(OFFSET_KEY, UNDEFINED_VALUE);
52-
long upstreamOffset = metadata.getOrDefault(UPSTREAM_OFFSET_KEY, UNDEFINED_VALUE);
53-
return new Pair<>(localOffset, upstreamOffset);
54-
}
55-
5623
public static Pair<PubSubPosition, PubSubPosition> getPositions(AdminMetadata metadata) {
5724
return new Pair<>(metadata.getPosition(), metadata.getUpstreamPosition());
5825
}
5926

60-
/**
61-
* @return the value to which the specified key is mapped to {@linkplain AdminTopicMetadataAccessor#EXECUTION_ID_KEY}.
62-
*/
63-
public static long getExecutionId(Map<String, Long> metadata) {
64-
return metadata.getOrDefault(EXECUTION_ID_KEY, UNDEFINED_VALUE);
65-
}
66-
6727
/**
6828
* @return the value to which the specified key is mapped to {@linkplain AdminTopicMetadataAccessor#ADMIN_OPERATION_PROTOCOL_VERSION_KEY}.
6929
*/
7030
public static long getAdminOperationProtocolVersion(AdminMetadata metadata) {
7131
return metadata.getAdminOperationProtocolVersion();
7232
}
7333

74-
/**
75-
* @return a pair of values representing local and upstream offsets
76-
*/
77-
public static Pair<Long, Long> getOffsets(AdminMetadata metadata) {
78-
long localOffset = metadata.getOffset() != null ? metadata.getOffset() : UNDEFINED_VALUE;
79-
long upstreamOffset = metadata.getUpstreamOffset() != null ? metadata.getUpstreamOffset() : UNDEFINED_VALUE;
80-
return new Pair<>(localOffset, upstreamOffset);
81-
}
82-
8334
/**
8435
* @return the execution ID from the metadata
8536
*/

0 commit comments

Comments
 (0)