From b0446c9caabef369e5186f89079ec1c1e84013dd Mon Sep 17 00:00:00 2001 From: Yang Guo Date: Mon, 27 Jul 2026 18:03:56 +0800 Subject: [PATCH] feat: return isr metadata --- fluss-rpc/src/main/proto/FlussApi.proto | 6 +- .../coordinator/CoordinatorRequestBatch.java | 14 +- .../fluss/server/metadata/BucketMetadata.java | 44 +++++- .../metadata/CoordinatorMetadataProvider.java | 39 ++++-- .../server/utils/ServerRpcMessageUtils.java | 13 +- .../fluss/server/zk/ZooKeeperClient.java | 7 +- .../CoordinatorEventProcessorTest.java | 14 +- .../server/metadata/BucketMetadataTest.java | 83 ++++++++++++ .../CoordinatorMetadataProviderTest.java | 69 ++++++++++ .../metadata/ZkBasedMetadataProviderTest.java | 19 ++- .../utils/ServerRpcMessageUtilsTest.java | 126 ++++++++++++++++++ 11 files changed, 411 insertions(+), 23 deletions(-) create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/metadata/BucketMetadataTest.java create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/metadata/CoordinatorMetadataProviderTest.java create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/utils/ServerRpcMessageUtilsTest.java diff --git a/fluss-rpc/src/main/proto/FlussApi.proto b/fluss-rpc/src/main/proto/FlussApi.proto index 57e96092d4..deb49fce30 100644 --- a/fluss-rpc/src/main/proto/FlussApi.proto +++ b/fluss-rpc/src/main/proto/FlussApi.proto @@ -869,8 +869,10 @@ message PbBucketMetadata { // optional as some time the leader may not elected yet optional int32 leader_id = 2; repeated int32 replica_id = 3 [packed = true]; - // TODO: Add isr here. optional int32 leader_epoch = 4; + // Generation of the complete leader/ISR state. Absence indicates legacy metadata. + optional int32 bucket_epoch = 5; + repeated int32 isr_id = 6 [packed = true]; } message PbProduceLogReqForBucket { @@ -1410,4 +1412,4 @@ message PbLiteralValue { optional int64 timestamp_millis_value = 11; // Epoch millis optional int32 timestamp_nano_of_millis_value = 12; // Nano of millis optional bytes decimal_bytes = 13; // Serialized decimal (non-compact mode) -} \ No newline at end of file +} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorRequestBatch.java b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorRequestBatch.java index 3ec69feeb6..bd5fa6b9f8 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorRequestBatch.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/coordinator/CoordinatorRequestBatch.java @@ -324,6 +324,12 @@ public void addUpdateMetadataRequestForTabletServers( Integer leaderEpoch = bucketLeaderAndIsr.map(LeaderAndIsr::leaderEpoch).orElse(null); Integer leader = bucketLeaderAndIsr.map(LeaderAndIsr::leader).orElse(null); + List isr = + bucketLeaderAndIsr.map(LeaderAndIsr::isr).orElse(Collections.emptyList()); + int bucketEpoch = + bucketLeaderAndIsr + .map(LeaderAndIsr::bucketEpoch) + .orElse(BucketMetadata.NO_LEADER_ISR_STATE_EPOCH); if (currentPartitionId == null) { Map> tableAssignment = coordinatorContext.getTableAssignment(currentTableId); @@ -332,7 +338,9 @@ public void addUpdateMetadataRequestForTabletServers( tableBucket.getBucket(), leader, leaderEpoch, - tableAssignment.get(tableBucket.getBucket())); + tableAssignment.get(tableBucket.getBucket()), + isr, + bucketEpoch); updateMetadataRequestBucketMap .computeIfAbsent(currentTableId, k -> new ArrayList<>()) .add(bucketMetadata); @@ -346,7 +354,9 @@ public void addUpdateMetadataRequestForTabletServers( tableBucket.getBucket(), leader, leaderEpoch, - partitionAssignment.get(tableBucket.getBucket())); + partitionAssignment.get(tableBucket.getBucket()), + isr, + bucketEpoch); updateMetadataRequestPartitionMap .computeIfAbsent(currentTableId, k -> new HashMap<>()) .computeIfAbsent(tableBucket.getPartitionId(), k -> new ArrayList<>()) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metadata/BucketMetadata.java b/fluss-server/src/main/java/org/apache/fluss/server/metadata/BucketMetadata.java index b6b968f28e..852b73c098 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metadata/BucketMetadata.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metadata/BucketMetadata.java @@ -19,26 +19,50 @@ import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Objects; import java.util.OptionalInt; /** This entity used to describe the bucket metadata. */ public class BucketMetadata { + public static final int NO_LEADER_ISR_STATE_EPOCH = -1; + private final int bucketId; private final @Nullable Integer leaderId; private final @Nullable Integer leaderEpoch; private final List replicas; + private final List isr; + private final @Nullable Integer bucketEpoch; + /** + * Creates legacy bucket metadata without authoritative leader/ISR state. + * + *

The empty ISR and absent bucket epoch distinguish this metadata from an authoritative + * state with an empty ISR. + */ public BucketMetadata( int bucketId, @Nullable Integer leaderId, @Nullable Integer leaderEpoch, List replicas) { + this(bucketId, leaderId, leaderEpoch, replicas, Collections.emptyList(), null); + } + + public BucketMetadata( + int bucketId, + @Nullable Integer leaderId, + @Nullable Integer leaderEpoch, + List replicas, + List isr, + @Nullable Integer bucketEpoch) { this.bucketId = bucketId; this.leaderId = leaderId; this.leaderEpoch = leaderEpoch; - this.replicas = replicas; + this.replicas = Collections.unmodifiableList(new ArrayList<>(replicas)); + this.isr = Collections.unmodifiableList(new ArrayList<>(isr)); + this.bucketEpoch = bucketEpoch; } public int getBucketId() { @@ -57,6 +81,14 @@ public List getReplicas() { return replicas; } + public List getIsr() { + return isr; + } + + public @Nullable Integer getBucketEpoch() { + return bucketEpoch; + } + @Override public String toString() { return "BucketMetadata{" @@ -68,6 +100,10 @@ public String toString() { + leaderEpoch + ", replicas=" + replicas + + ", isr=" + + isr + + ", bucketEpoch=" + + bucketEpoch + '}'; } @@ -83,11 +119,13 @@ public boolean equals(Object o) { return bucketId == that.bucketId && Objects.equals(leaderId, that.leaderId) && Objects.equals(leaderEpoch, that.leaderEpoch) - && replicas.equals(that.replicas); + && replicas.equals(that.replicas) + && isr.equals(that.isr) + && Objects.equals(bucketEpoch, that.bucketEpoch); } @Override public int hashCode() { - return Objects.hash(bucketId, leaderId, leaderEpoch, replicas); + return Objects.hash(bucketId, leaderId, leaderEpoch, replicas, isr, bucketEpoch); } } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metadata/CoordinatorMetadataProvider.java b/fluss-server/src/main/java/org/apache/fluss/server/metadata/CoordinatorMetadataProvider.java index 27511d0540..9cae3977f7 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metadata/CoordinatorMetadataProvider.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metadata/CoordinatorMetadataProvider.java @@ -17,6 +17,7 @@ package org.apache.fluss.server.metadata; +import org.apache.fluss.annotation.VisibleForTesting; import org.apache.fluss.metadata.PhysicalTablePath; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TableInfo; @@ -30,6 +31,7 @@ import javax.annotation.Nullable; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Optional; @@ -115,9 +117,9 @@ public Optional getPhysicalTablePathFromCache(long partitionI * Constructs bucket metadata list from coordinator context information. * *

This method builds a complete list of bucket metadata by combining assignment information - * from the provided table assignment map with leader and epoch data from the coordinator - * context. Each bucket's metadata includes its ID, current leader server, leader epoch, and - * replica assignments. + * from the provided table assignment map with leader and ISR data from the coordinator context. + * Each bucket's metadata includes its ID, current leader server, leader epoch, ISR, bucket + * epoch, and replica assignments. * * @param ctx the coordinator context containing leader and epoch information * @param tableId the table identifier @@ -125,7 +127,8 @@ public Optional getPhysicalTablePathFromCache(long partitionI * @param tableAssignment the assignment map from bucket ID to list of replica server IDs * @return a list of bucket metadata objects containing complete bucket information */ - private static List getBucketMetadataFromContext( + @VisibleForTesting + static List getBucketMetadataFromContext( CoordinatorContext ctx, long tableId, @Nullable Long partitionId, @@ -135,13 +138,27 @@ private static List getBucketMetadataFromContext( (bucketId, serverIds) -> { TableBucket tableBucket = new TableBucket(tableId, partitionId, bucketId); Optional optLeaderAndIsr = ctx.getBucketLeaderAndIsr(tableBucket); - Integer leader = optLeaderAndIsr.map(LeaderAndIsr::leader).orElse(null); - BucketMetadata bucketMetadata = - new BucketMetadata( - bucketId, - leader, - ctx.getBucketLeaderEpoch(tableBucket), - serverIds); + BucketMetadata bucketMetadata; + if (optLeaderAndIsr.isPresent()) { + LeaderAndIsr leaderAndIsr = optLeaderAndIsr.get(); + bucketMetadata = + new BucketMetadata( + bucketId, + leaderAndIsr.leader(), + leaderAndIsr.leaderEpoch(), + serverIds, + leaderAndIsr.isr(), + leaderAndIsr.bucketEpoch()); + } else { + bucketMetadata = + new BucketMetadata( + bucketId, + null, + null, + serverIds, + Collections.emptyList(), + BucketMetadata.NO_LEADER_ISR_STATE_EPOCH); + } bucketMetadataList.add(bucketMetadata); }); return bucketMetadataList; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/utils/ServerRpcMessageUtils.java b/fluss-server/src/main/java/org/apache/fluss/server/utils/ServerRpcMessageUtils.java index 1e5e914c0d..efff82bfa5 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/utils/ServerRpcMessageUtils.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/utils/ServerRpcMessageUtils.java @@ -648,6 +648,15 @@ private static List toPbBucketMetadata( pbBucketMetadata.addReplicaId(replica); } + Integer bucketEpoch = bucketMetadata.getBucketEpoch(); + if (bucketEpoch != null) { + pbBucketMetadata.setBucketEpoch(bucketEpoch); + } + + for (Integer isr : bucketMetadata.getIsr()) { + pbBucketMetadata.addIsrId(isr); + } + pbBucketMetadataList.add(pbBucketMetadata); } return pbBucketMetadataList; @@ -688,7 +697,9 @@ private static BucketMetadata toBucketMetadata(PbBucketMetadata pbBucketMetadata pbBucketMetadata.hasLeaderEpoch() ? pbBucketMetadata.getLeaderEpoch() : null, Arrays.stream(pbBucketMetadata.getReplicaIds()) .boxed() - .collect(Collectors.toList())); + .collect(Collectors.toList()), + Arrays.stream(pbBucketMetadata.getIsrIds()).boxed().collect(Collectors.toList()), + pbBucketMetadata.hasBucketEpoch() ? pbBucketMetadata.getBucketEpoch() : null); } private static PartitionMetadata toPartitionMetadata(PbPartitionMetadata pbPartitionMetadata) { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/zk/ZooKeeperClient.java b/fluss-server/src/main/java/org/apache/fluss/server/zk/ZooKeeperClient.java index e9eab00e00..7a900897f1 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/zk/ZooKeeperClient.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/zk/ZooKeeperClient.java @@ -1780,8 +1780,13 @@ private BucketMetadata createBucketMetadata( LeaderAndIsr leaderAndIsr = leaderAndIsrs.get(bucket); Integer leader = leaderAndIsr != null ? leaderAndIsr.leader() : null; Integer leaderEpoch = leaderAndIsr != null ? leaderAndIsr.leaderEpoch() : null; + List isr = leaderAndIsr != null ? leaderAndIsr.isr() : Collections.emptyList(); + int bucketEpoch = + leaderAndIsr != null + ? leaderAndIsr.bucketEpoch() + : BucketMetadata.NO_LEADER_ISR_STATE_EPOCH; List replicas = assignment.getBucketAssignments().get(bucketId).getReplicas(); - return new BucketMetadata(bucketId, leader, leaderEpoch, replicas); + return new BucketMetadata(bucketId, leader, leaderEpoch, replicas, isr, bucketEpoch); } /** Close the underlying ZooKeeperClient. */ diff --git a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessorTest.java b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessorTest.java index 5a5a68b3c8..ddd11160c4 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessorTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/coordinator/CoordinatorEventProcessorTest.java @@ -1532,7 +1532,12 @@ void testSchemaChange() throws Exception { tableInfo, Collections.singletonList( new BucketMetadata( - 0, replicas.get(0), 0, replicas))))); + 0, + replicas.get(0), + 0, + replicas, + replicas, + 0))))); // alter table column. alterTable( @@ -1591,7 +1596,12 @@ void testTableRegistrationChange() throws Exception { tableInfo, Collections.singletonList( new BucketMetadata( - 0, replicas.get(0), 0, replicas))))); + 0, + replicas.get(0), + 0, + replicas, + replicas, + 0))))); // alter table properties (custom property) TablePropertyChanges.Builder builder = TablePropertyChanges.builder(); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/metadata/BucketMetadataTest.java b/fluss-server/src/test/java/org/apache/fluss/server/metadata/BucketMetadataTest.java new file mode 100644 index 0000000000..281a0e4a50 --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/metadata/BucketMetadataTest.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.metadata; + +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class BucketMetadataTest { + + @Test + void testDefensivelyCopiesAndExposesImmutableReplicaAndIsrLists() { + List replicas = new ArrayList<>(Arrays.asList(1, 2, 3)); + List isr = new ArrayList<>(Arrays.asList(1, 2)); + BucketMetadata metadata = new BucketMetadata(0, 1, 2, replicas, isr, 3); + + replicas.add(4); + isr.add(3); + + assertThat(metadata.getReplicas()).containsExactly(1, 2, 3); + assertThat(metadata.getIsr()).containsExactly(1, 2); + assertThatThrownBy(() -> metadata.getReplicas().add(4)) + .isInstanceOf(UnsupportedOperationException.class); + assertThatThrownBy(() -> metadata.getIsr().add(3)) + .isInstanceOf(UnsupportedOperationException.class); + } + + @Test + void testEqualityHashCodeAndStringIncludeAuthoritativeState() { + BucketMetadata metadata = + new BucketMetadata(0, 1, 2, Arrays.asList(1, 2, 3), Arrays.asList(1, 2), 3); + BucketMetadata sameMetadata = + new BucketMetadata(0, 1, 2, Arrays.asList(1, 2, 3), Arrays.asList(1, 2), 3); + BucketMetadata differentIsr = + new BucketMetadata( + 0, 1, 2, Arrays.asList(1, 2, 3), Collections.singletonList(1), 3); + BucketMetadata differentEpoch = + new BucketMetadata(0, 1, 2, Arrays.asList(1, 2, 3), Arrays.asList(1, 2), 4); + + assertThat(metadata).isEqualTo(sameMetadata).hasSameHashCodeAs(sameMetadata); + assertThat(metadata).isNotEqualTo(differentIsr).isNotEqualTo(differentEpoch); + assertThat(metadata.toString()).contains("isr=[1, 2]", "bucketEpoch=3"); + } + + @Test + void testAuthoritativeEmptyIsrDiffersFromLegacyUnknownIsr() { + BucketMetadata authoritative = + new BucketMetadata( + 0, + null, + null, + Arrays.asList(1, 2, 3), + Collections.emptyList(), + BucketMetadata.NO_LEADER_ISR_STATE_EPOCH); + BucketMetadata legacy = new BucketMetadata(0, null, null, Arrays.asList(1, 2, 3)); + + assertThat(authoritative.getBucketEpoch()) + .isEqualTo(BucketMetadata.NO_LEADER_ISR_STATE_EPOCH); + assertThat(legacy.getBucketEpoch()).isNull(); + assertThat(authoritative).isNotEqualTo(legacy); + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/metadata/CoordinatorMetadataProviderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/metadata/CoordinatorMetadataProviderTest.java new file mode 100644 index 0000000000..9ab0a8a80c --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/metadata/CoordinatorMetadataProviderTest.java @@ -0,0 +1,69 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.metadata; + +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.server.coordinator.CoordinatorContext; +import org.apache.fluss.server.zk.ZkEpoch; +import org.apache.fluss.server.zk.data.LeaderAndIsr; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +class CoordinatorMetadataProviderTest { + + @Test + void testBuildsAuthoritativeBucketMetadataFromSingleLeaderAndIsrState() { + long tableId = 10L; + CoordinatorContext context = new CoordinatorContext(ZkEpoch.INITIAL_EPOCH); + LeaderAndIsr leaderAndIsr = + new LeaderAndIsr(1, 2, Arrays.asList(1, 2), Collections.emptyList(), 3, 4); + context.putBucketLeaderAndIsr(new TableBucket(tableId, 0), leaderAndIsr); + Map> assignment = new LinkedHashMap<>(); + assignment.put(0, Arrays.asList(1, 2, 3)); + assignment.put(1, Arrays.asList(2, 3, 4)); + + List metadata = + CoordinatorMetadataProvider.getBucketMetadataFromContext( + context, tableId, null, assignment); + + assertThat(metadata) + .containsExactly( + new BucketMetadata( + 0, + leaderAndIsr.leader(), + leaderAndIsr.leaderEpoch(), + Arrays.asList(1, 2, 3), + leaderAndIsr.isr(), + leaderAndIsr.bucketEpoch()), + new BucketMetadata( + 1, + null, + null, + Arrays.asList(2, 3, 4), + Collections.emptyList(), + BucketMetadata.NO_LEADER_ISR_STATE_EPOCH)); + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/metadata/ZkBasedMetadataProviderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/metadata/ZkBasedMetadataProviderTest.java index a2f4cbedb9..5e74623c74 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/metadata/ZkBasedMetadataProviderTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/metadata/ZkBasedMetadataProviderTest.java @@ -100,6 +100,7 @@ void testGetTableMetadataFromZk() throws Exception { TableAssignment.builder() .add(0, BucketAssignment.of(1, 2, 3)) .add(1, BucketAssignment.of(2, 3, 4)) + .add(2, BucketAssignment.of(3, 4, 5)) .build(); metadataManager.createDatabase("test_db", DatabaseDescriptor.EMPTY, true); long tableId = @@ -132,7 +133,7 @@ void testGetTableMetadataFromZk() throws Exception { assertThat(tableInfo.getSchema()).isEqualTo(desc.getSchema()); List bucketMetadataList = tableMetadata.getBucketMetadataList(); - assertThat(bucketMetadataList).hasSize(2); + assertThat(bucketMetadataList).hasSize(3); // Verify bucket metadata Map bucketMap = @@ -142,10 +143,22 @@ void testGetTableMetadataFromZk() throws Exception { BucketMetadata bucket0Metadata = bucketMap.get(0); assertThat(extractLeaderFromBucketMetadata(bucket0Metadata)).isEqualTo(1); assertThat(bucket0Metadata.getReplicas()).containsExactly(1, 2, 3); + assertThat(bucket0Metadata.getIsr()).containsExactly(1, 2, 3); + assertThat(bucket0Metadata.getBucketEpoch()).isEqualTo(1000); BucketMetadata bucket1Metadata = bucketMap.get(1); assertThat(extractLeaderFromBucketMetadata(bucket1Metadata)).isEqualTo(2); assertThat(bucket1Metadata.getReplicas()).containsExactly(2, 3, 4); + assertThat(bucket1Metadata.getIsr()).containsExactly(2, 3, 4); + assertThat(bucket1Metadata.getBucketEpoch()).isEqualTo(2000); + + BucketMetadata bucket2Metadata = bucketMap.get(2); + assertThat(bucket2Metadata.getLeaderId()).isEmpty(); + assertThat(bucket2Metadata.getLeaderEpoch()).isEmpty(); + assertThat(bucket2Metadata.getReplicas()).containsExactly(3, 4, 5); + assertThat(bucket2Metadata.getIsr()).isEmpty(); + assertThat(bucket2Metadata.getBucketEpoch()) + .isEqualTo(BucketMetadata.NO_LEADER_ISR_STATE_EPOCH); } @Test @@ -213,10 +226,14 @@ void testGetPartitionMetadataFromZk() throws Exception { BucketMetadata bucket0Metadata = bucketMap.get(0); assertThat(extractLeaderFromBucketMetadata(bucket0Metadata)).isEqualTo(1); assertThat(bucket0Metadata.getReplicas()).containsExactly(1, 2); + assertThat(bucket0Metadata.getIsr()).containsExactly(1, 2); + assertThat(bucket0Metadata.getBucketEpoch()).isEqualTo(1000); BucketMetadata bucket1Metadata = bucketMap.get(1); assertThat(extractLeaderFromBucketMetadata(bucket1Metadata)).isEqualTo(2); assertThat(bucket1Metadata.getReplicas()).containsExactly(2, 3); + assertThat(bucket1Metadata.getIsr()).containsExactly(2, 3); + assertThat(bucket1Metadata.getBucketEpoch()).isEqualTo(2000); } @Test diff --git a/fluss-server/src/test/java/org/apache/fluss/server/utils/ServerRpcMessageUtilsTest.java b/fluss-server/src/test/java/org/apache/fluss/server/utils/ServerRpcMessageUtilsTest.java new file mode 100644 index 0000000000..49f7840e4c --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/utils/ServerRpcMessageUtilsTest.java @@ -0,0 +1,126 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.utils; + +import org.apache.fluss.rpc.messages.PbBucketMetadata; +import org.apache.fluss.rpc.messages.PbPartitionMetadata; +import org.apache.fluss.rpc.messages.PbTableMetadata; +import org.apache.fluss.rpc.messages.UpdateMetadataRequest; +import org.apache.fluss.server.metadata.BucketMetadata; +import org.apache.fluss.server.metadata.ClusterMetadata; +import org.apache.fluss.server.metadata.PartitionMetadata; +import org.apache.fluss.server.metadata.TableMetadata; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.apache.fluss.record.TestData.DATA1_TABLE_INFO; +import static org.apache.fluss.server.metadata.BucketMetadata.NO_LEADER_ISR_STATE_EPOCH; +import static org.apache.fluss.server.utils.ServerRpcMessageUtils.getUpdateMetadataRequestData; +import static org.apache.fluss.server.utils.ServerRpcMessageUtils.makeUpdateMetadataRequest; +import static org.assertj.core.api.Assertions.assertThat; + +class ServerRpcMessageUtilsTest { + + @Test + void testAuthoritativeBucketMetadataTableAndPartitionRoundTrip() { + List bucketMetadata = + Arrays.asList( + new BucketMetadata(0, 1, 2, Arrays.asList(1, 2, 3), Arrays.asList(1, 2), 3), + new BucketMetadata( + 1, + null, + null, + Arrays.asList(1, 2, 3), + Collections.emptyList(), + NO_LEADER_ISR_STATE_EPOCH), + new BucketMetadata( + 2, 2, 4, Arrays.asList(2, 3), Collections.emptyList(), 5)); + TableMetadata tableMetadata = new TableMetadata(DATA1_TABLE_INFO, bucketMetadata); + PartitionMetadata partitionMetadata = + new PartitionMetadata(DATA1_TABLE_INFO.getTableId(), "p1", 10L, bucketMetadata); + + UpdateMetadataRequest request = + makeUpdateMetadataRequest( + null, + null, + Collections.emptySet(), + Collections.singletonList(tableMetadata), + Collections.singletonList(partitionMetadata)); + + PbTableMetadata pbTableMetadata = request.getTableMetadatasList().get(0); + assertAuthoritativeWireMetadata(pbTableMetadata.getBucketMetadatasList()); + PbPartitionMetadata pbPartitionMetadata = request.getPartitionMetadatasList().get(0); + assertAuthoritativeWireMetadata(pbPartitionMetadata.getBucketMetadatasList()); + + ClusterMetadata decoded = getUpdateMetadataRequestData(request); + assertThat(decoded.getTableMetadataList().get(0).getBucketMetadataList()) + .containsExactlyElementsOf(bucketMetadata); + assertThat(decoded.getPartitionMetadataList().get(0).getBucketMetadataList()) + .containsExactlyElementsOf(bucketMetadata); + } + + @Test + void testLegacyBucketMetadataRoundTripOmitsBucketEpoch() { + BucketMetadata legacy = new BucketMetadata(0, 1, 2, Arrays.asList(1, 2, 3)); + UpdateMetadataRequest request = + makeUpdateMetadataRequest( + null, + null, + Collections.emptySet(), + Collections.singletonList( + new TableMetadata( + DATA1_TABLE_INFO, Collections.singletonList(legacy))), + Collections.emptyList()); + + PbBucketMetadata pbBucketMetadata = + request.getTableMetadatasList().get(0).getBucketMetadatasList().get(0); + assertThat(pbBucketMetadata.hasBucketEpoch()).isFalse(); + assertThat(pbBucketMetadata.getIsrIds()).isEmpty(); + + BucketMetadata decoded = + getUpdateMetadataRequestData(request) + .getTableMetadataList() + .get(0) + .getBucketMetadataList() + .get(0); + assertThat(decoded).isEqualTo(legacy); + assertThat(decoded.getBucketEpoch()).isNull(); + } + + private static void assertAuthoritativeWireMetadata(List bucketMetadata) { + assertThat(bucketMetadata.get(0).getBucketId()).isEqualTo(0); + assertThat(bucketMetadata.get(0).getLeaderId()).isEqualTo(1); + assertThat(bucketMetadata.get(0).getLeaderEpoch()).isEqualTo(2); + assertThat(bucketMetadata.get(0).getReplicaIds()).containsExactly(1, 2, 3); + assertThat(bucketMetadata.get(0).hasBucketEpoch()).isTrue(); + assertThat(bucketMetadata.get(0).getBucketEpoch()).isEqualTo(3); + assertThat(bucketMetadata.get(0).getIsrIds()).containsExactly(1, 2); + + assertThat(bucketMetadata.get(1).hasBucketEpoch()).isTrue(); + assertThat(bucketMetadata.get(1).getBucketEpoch()).isEqualTo(NO_LEADER_ISR_STATE_EPOCH); + assertThat(bucketMetadata.get(1).getIsrIds()).isEmpty(); + + assertThat(bucketMetadata.get(2).hasBucketEpoch()).isTrue(); + assertThat(bucketMetadata.get(2).getBucketEpoch()).isEqualTo(5); + assertThat(bucketMetadata.get(2).getIsrIds()).isEmpty(); + } +}