This is an automated email from the ASF dual-hosted git repository.
rpuch pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/ignite-3.git
The following commit(s) were added to refs/heads/main by this push:
new b7782257f22 IGNITE-25678 Fix
laggingSchemasPreventPartitionDataReplication (#6045)
b7782257f22 is described below
commit b7782257f2299884b45fdea6f2b29edb3b49b2be
Author: Roman Puchkovskiy <[email protected]>
AuthorDate: Mon Jun 16 11:38:58 2025 +0400
IGNITE-25678 Fix laggingSchemasPreventPartitionDataReplication (#6045)
---
.../raftsnapshot/ItTableRaftSnapshotsTest.java | 28 +---------------
.../schemasync/ItSchemaSyncAndReplicationTest.java | 6 +++-
.../java/org/apache/ignite/internal/Cluster.java | 39 ++++++++++++++++++++--
3 files changed, 43 insertions(+), 30 deletions(-)
diff --git
a/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
b/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
index ba2eeb16bce..b6702745ca2 100644
---
a/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
+++
b/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
@@ -59,7 +59,6 @@ import java.util.stream.Stream;
import org.apache.ignite.Ignite;
import org.apache.ignite.internal.Cluster;
import org.apache.ignite.internal.ClusterPerTestIntegrationTest;
-import org.apache.ignite.internal.TestWrappers;
import org.apache.ignite.internal.app.IgniteImpl;
import org.apache.ignite.internal.cluster.management.CmgGroupId;
import org.apache.ignite.internal.logger.IgniteLogger;
@@ -80,7 +79,6 @@ import
org.apache.ignite.internal.storage.pagememory.PersistentPageMemoryStorage
import
org.apache.ignite.internal.storage.pagememory.VolatilePageMemoryStorageEngine;
import org.apache.ignite.internal.storage.rocksdb.RocksDbStorageEngine;
import org.apache.ignite.internal.table.InternalTable;
-import org.apache.ignite.internal.table.NodeUtils;
import
org.apache.ignite.internal.table.distributed.schema.PartitionCommandsMarshallerImpl;
import org.apache.ignite.internal.testframework.IgniteTestUtils;
import org.apache.ignite.internal.testframework.log4j2.LogInspector;
@@ -325,19 +323,6 @@ class ItTableRaftSnapshotsTest extends
ClusterPerTestIntegrationTest {
LOG.info("Lease is accepted by [nodeConsistentId={}].",
primary.join().getLeaseholder());
}
- private @Nullable String getPrimaryReplicaName() {
- IgniteImpl node = unwrapIgniteImpl(cluster.node(0));
-
- CompletableFuture<ReplicaMeta> primary =
node.placementDriver().getPrimaryReplica(
- cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME),
- node.clockService().now()
- );
-
- assertThat(primary, willCompleteSuccessfully());
-
- return primary.join().getLeaseholder();
- }
-
private void putToNode(int nodeIndex, int key, String value) {
putToNode(nodeIndex, key, value, null);
}
@@ -456,18 +441,7 @@ class ItTableRaftSnapshotsTest extends
ClusterPerTestIntegrationTest {
}
private void transferPrimaryOnSolePartitionTo(int nodeIndex) throws
InterruptedException {
- String proposedPrimaryName = cluster.node(nodeIndex).name();
-
- if (!proposedPrimaryName.equals(getPrimaryReplicaName())) {
-
- String newPrimaryName = NodeUtils.transferPrimary(
-
cluster.runningNodes().map(TestWrappers::unwrapIgniteImpl).collect(toList()),
- cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME),
- proposedPrimaryName
- );
-
- assertEquals(proposedPrimaryName, newPrimaryName);
- }
+ cluster.transferPrimaryTo(nodeIndex,
cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME));
}
/**
diff --git
a/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
b/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
index 449ec1a32f1..c216102f31a 100644
---
a/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
+++
b/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
@@ -35,6 +35,7 @@ import
org.apache.ignite.internal.ClusterPerTestIntegrationTest;
import org.apache.ignite.internal.app.IgniteImpl;
import org.apache.ignite.internal.metastorage.server.WatchListenerInhibitor;
import org.apache.ignite.internal.metastorage.server.raft.MetastorageGroupId;
+import org.apache.ignite.internal.replicator.ReplicationGroupId;
import org.apache.ignite.internal.storage.MvPartitionStorage;
import org.apache.ignite.internal.storage.RowId;
import org.apache.ignite.internal.table.TableViewInternal;
@@ -118,7 +119,10 @@ class ItSchemaSyncAndReplicationTest extends
ClusterPerTestIntegrationTest {
private void transferLeadershipsTo(int nodeIndex) throws
InterruptedException {
cluster.transferLeadershipTo(nodeIndex, MetastorageGroupId.INSTANCE);
- cluster.transferLeadershipTo(nodeIndex,
cluster.solePartitionId(ZONE_NAME, TABLE_NAME));
+
+ ReplicationGroupId solePartitionId =
cluster.solePartitionId(ZONE_NAME, TABLE_NAME);
+ cluster.transferLeadershipTo(nodeIndex, solePartitionId);
+ cluster.transferPrimaryTo(nodeIndex, solePartitionId);
}
private CompletableFuture<?> rejectionDueToMetadataLagTriggered() {
diff --git
a/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
b/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
index 52b0082569b..6350451d9a6 100644
---
a/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
+++
b/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
@@ -18,6 +18,7 @@
package org.apache.ignite.internal;
import static java.util.Collections.nCopies;
+import static java.util.concurrent.TimeUnit.SECONDS;
import static java.util.stream.Collectors.joining;
import static java.util.stream.Collectors.toList;
import static org.apache.ignite.internal.ClusterConfiguration.configOverrides;
@@ -47,7 +48,6 @@ import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
-import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BiPredicate;
import java.util.function.BooleanSupplier;
@@ -68,10 +68,12 @@ import
org.apache.ignite.internal.lang.IgniteStringFormatter;
import org.apache.ignite.internal.logger.IgniteLogger;
import org.apache.ignite.internal.logger.Loggers;
import org.apache.ignite.internal.network.NetworkMessage;
+import org.apache.ignite.internal.placementdriver.ReplicaMeta;
import org.apache.ignite.internal.raft.RaftNodeId;
import org.apache.ignite.internal.raft.server.impl.JraftServerImpl;
import org.apache.ignite.internal.replicator.ReplicationGroupId;
import org.apache.ignite.internal.sql.SqlCommon;
+import org.apache.ignite.internal.table.NodeUtils;
import org.apache.ignite.internal.testframework.TestIgnitionManager;
import org.apache.ignite.raft.jraft.RaftGroupService;
import org.apache.ignite.raft.jraft.Status;
@@ -453,7 +455,7 @@ public class Cluster {
*/
public Ignite startNode(int index, String nodeBootstrapConfigTemplate) {
ServerRegistration registration = startEmbeddedNode(index,
nodeBootstrapConfigTemplate);
- assertThat("nodeIndex=" + index, registration.registrationFuture(),
willSucceedIn(20, TimeUnit.SECONDS));
+ assertThat("nodeIndex=" + index, registration.registrationFuture(),
willSucceedIn(20, SECONDS));
Ignite newIgniteNode = registration.server().api();
assertEquals(newIgniteNode, nodes.get(index));
@@ -755,6 +757,39 @@ public class Cluster {
}
}
+ /**
+ * Transfers primary replica of given replication group to the node with
given index.
+ *
+ * @param nodeIndex Destination node index.
+ * @param groupId ID of the replication group.
+ */
+ public void transferPrimaryTo(int nodeIndex, ReplicationGroupId groupId)
throws InterruptedException {
+ String proposedPrimaryName = node(nodeIndex).name();
+
+ if (!proposedPrimaryName.equals(getPrimaryReplicaName(groupId))) {
+
+ String newPrimaryName = NodeUtils.transferPrimary(
+
runningNodes().map(TestWrappers::unwrapIgniteImpl).collect(toList()),
+ groupId,
+ proposedPrimaryName
+ );
+
+ assertEquals(proposedPrimaryName, newPrimaryName);
+ }
+ }
+
+ private @Nullable String getPrimaryReplicaName(ReplicationGroupId groupId)
{
+ IgniteImpl node = unwrapIgniteImpl(aliveNode());
+
+ CompletableFuture<ReplicaMeta> primary = node.placementDriver()
+ .awaitPrimaryReplica(groupId, node.clockService().now(), 30,
SECONDS);
+
+ assertThat(primary, willCompleteSuccessfully());
+
+ @Nullable ReplicaMeta replicaMeta = primary.join();
+ return replicaMeta != null ? replicaMeta.getLeaseholder() : null;
+ }
+
/**
* Returns the ID of the sole partition that exists in the cluster or
throws if there are less than one
* or more than one partitions.