diff --git a/CHANGES.txt b/CHANGES.txt index 2ba47bb462..ff1418e1d5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.1 + * Reconfigure CMS before assassinate (CASSANDRA-19768) * Warn about unqualified prepared statement only if it is select or modification statement (CASSANDRA-18322) * Update legacy peers tables during node replacement (CASSANDRA-19782) * Refactor ColumnCondition (CASSANDRA-19620) diff --git a/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java b/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java index 400553ce01..6ba47ae04d 100644 --- a/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java +++ b/src/java/org/apache/cassandra/locator/CMSPlacementStrategy.java @@ -48,7 +48,8 @@ import static org.apache.cassandra.locator.SimpleStrategy.REPLICATION_FACTOR; */ public interface CMSPlacementStrategy { - Set reconfigure(Set currentCms, ClusterMetadata metadata); + Set reconfigure(ClusterMetadata metadata); + boolean needsReconfiguration(ClusterMetadata metadata); static CMSPlacementStrategy fromReplicationParams(ReplicationParams params, Predicate filter) { @@ -92,20 +93,24 @@ public interface CMSPlacementStrategy this.filter = filter; } - public Set reconfigure(Set currentCms, ClusterMetadata metadata) + public Set reconfigure(ClusterMetadata metadata) { Map rf = new HashMap<>(this.rf.size()); for (Map.Entry e : this.rf.entrySet()) { Collection nodesInDc = metadata.directory.allDatacenterEndpoints().get(e.getKey()); - if (nodesInDc == null) + if (nodesInDc.isEmpty()) throw new IllegalStateException(String.format("There are no nodes in %s datacenter", e.getKey())); if (nodesInDc.size() < e.getValue()) throw new Transformation.RejectedTransformationException(String.format("There are not enough nodes in %s datacenter to satisfy replication factor", e.getKey())); rf.put(e.getKey(), ReplicationFactor.fullOnly(e.getValue())); } + return reconfigure(metadata, rf); + } + public Set reconfigure(ClusterMetadata metadata, Map rf) + { Directory tmpDirectory = metadata.directory; TokenMap tokenMap = metadata.tokenMap; for (NodeId peerId : metadata.directory.peerIds()) @@ -129,6 +134,25 @@ public interface CMSPlacementStrategy return endpoints.endpoints().stream().map(metadata.directory::peerId).collect(Collectors.toSet()); } + + public boolean needsReconfiguration(ClusterMetadata metadata) + { + Map rf = new HashMap<>(this.rf.size()); + for (Map.Entry e : this.rf.entrySet()) + { + Collection nodesInDc = metadata.directory.allDatacenterEndpoints().get(e.getKey()); + if (nodesInDc.size() < e.getValue()) + return true; + rf.put(e.getKey(), ReplicationFactor.fullOnly(e.getValue())); + } + + Set currentCms = metadata.fullCMSMembers() + .stream() + .map(metadata.directory::peerId) + .collect(Collectors.toSet()); + Set newCms = reconfigure(metadata, rf); + return !currentCms.equals(newCms); + } } class DefaultNodeFilter implements BiFunction diff --git a/src/java/org/apache/cassandra/metrics/TCMMetrics.java b/src/java/org/apache/cassandra/metrics/TCMMetrics.java index 29858ead80..e58a3b11cd 100644 --- a/src/java/org/apache/cassandra/metrics/TCMMetrics.java +++ b/src/java/org/apache/cassandra/metrics/TCMMetrics.java @@ -25,12 +25,13 @@ import com.codahale.metrics.Histogram; import com.codahale.metrics.Meter; import com.codahale.metrics.Timer; import org.apache.cassandra.gms.FailureDetector; +import org.apache.cassandra.locator.CMSPlacementStrategy; +import org.apache.cassandra.schema.ReplicationParams; import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.utils.FBUtilities; import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics; -import static org.apache.cassandra.tcm.transformations.cms.PrepareCMSReconfiguration.needsReconfiguration; public class TCMMetrics { @@ -95,7 +96,10 @@ public class TCMMetrics needsCMSReconfiguration = Metrics.register(factory.createMetricName("NeedsCMSReconfiguration"), () -> { ClusterMetadata metadata = ClusterMetadata.currentNullable(); - return metadata != null && needsReconfiguration(metadata) ? 1 : 0; + if (metadata == null) + return 0; + CMSPlacementStrategy placementStrategy = CMSPlacementStrategy.fromReplicationParams(ReplicationParams.meta(metadata), nodeId -> true); + return placementStrategy.needsReconfiguration(metadata) ? 1 : 0; }); fetchedPeerLogEntries = Metrics.histogram(factory.createMetricName("FetchedPeerLogEntries"), false); diff --git a/src/java/org/apache/cassandra/tcm/CMSOperations.java b/src/java/org/apache/cassandra/tcm/CMSOperations.java index 2211a2bb3d..30ee48b20d 100644 --- a/src/java/org/apache/cassandra/tcm/CMSOperations.java +++ b/src/java/org/apache/cassandra/tcm/CMSOperations.java @@ -31,6 +31,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.locator.CMSPlacementStrategy; import org.apache.cassandra.schema.ReplicationParams; import org.apache.cassandra.tcm.membership.NodeId; import org.apache.cassandra.tcm.membership.NodeState; @@ -41,7 +42,6 @@ import org.apache.cassandra.tcm.sequences.ReconfigureCMS; import org.apache.cassandra.tcm.serialization.Version; import org.apache.cassandra.tcm.transformations.Unregister; import org.apache.cassandra.tcm.transformations.cms.AdvanceCMSReconfiguration; -import org.apache.cassandra.tcm.transformations.cms.PrepareCMSReconfiguration; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.MBeanWrapper; @@ -142,7 +142,8 @@ public class CMSOperations implements CMSOperationsMBean ClusterMetadata metadata = ClusterMetadata.current(); String members = metadata.fullCMSMembers().stream().sorted().map(Object::toString).collect(Collectors.joining(",")); info.put(MEMBERS, members); - info.put(NEEDS_RECONFIGURATION, Boolean.toString(PrepareCMSReconfiguration.needsReconfiguration(metadata))); + CMSPlacementStrategy placementStrategy = CMSPlacementStrategy.fromReplicationParams(ReplicationParams.meta(metadata), nodeId -> true); + info.put(NEEDS_RECONFIGURATION, Boolean.toString(placementStrategy.needsReconfiguration(metadata))); info.put(IS_MEMBER, Boolean.toString(cms.isCurrentMember(FBUtilities.getBroadcastAddressAndPort()))); info.put(SERVICE_STATE, ClusterMetadataService.state(metadata).toString()); info.put(IS_MIGRATING, Boolean.toString(cms.isMigrating())); diff --git a/src/java/org/apache/cassandra/tcm/ownership/TokenMap.java b/src/java/org/apache/cassandra/tcm/ownership/TokenMap.java index d8dca5ff52..c32f6c351c 100644 --- a/src/java/org/apache/cassandra/tcm/ownership/TokenMap.java +++ b/src/java/org/apache/cassandra/tcm/ownership/TokenMap.java @@ -135,8 +135,6 @@ public class TokenMap implements MetadataValue public ImmutableList tokens(NodeId nodeId) { Collection tokens = map.inverse().get(nodeId); - if (tokens == null) - return null; return ImmutableList.copyOf(tokens); } diff --git a/src/java/org/apache/cassandra/tcm/transformations/Assassinate.java b/src/java/org/apache/cassandra/tcm/transformations/Assassinate.java index 0cee0f627f..76bca253ee 100644 --- a/src/java/org/apache/cassandra/tcm/transformations/Assassinate.java +++ b/src/java/org/apache/cassandra/tcm/transformations/Assassinate.java @@ -30,6 +30,7 @@ import org.apache.cassandra.tcm.ownership.PlacementDeltas; import org.apache.cassandra.tcm.ownership.PlacementProvider; import org.apache.cassandra.tcm.sequences.LeaveStreams; import org.apache.cassandra.tcm.sequences.LockedRanges; +import org.apache.cassandra.tcm.sequences.ReconfigureCMS; import org.apache.cassandra.tcm.sequences.UnbootstrapAndLeave; public class Assassinate extends PrepareLeave @@ -60,6 +61,8 @@ public class Assassinate extends PrepareLeave if (!metadata.directory.isRegistered(endpoint)) return; + ReconfigureCMS.maybeReconfigureCMS(metadata, endpoint); + NodeId nodeId = metadata.directory.peerId(endpoint); ClusterMetadataService.instance().commit(new Assassinate(nodeId, ClusterMetadataService.instance().placementProvider())); diff --git a/src/java/org/apache/cassandra/tcm/transformations/cms/PrepareCMSReconfiguration.java b/src/java/org/apache/cassandra/tcm/transformations/cms/PrepareCMSReconfiguration.java index 0367640fda..2d41fb53a6 100644 --- a/src/java/org/apache/cassandra/tcm/transformations/cms/PrepareCMSReconfiguration.java +++ b/src/java/org/apache/cassandra/tcm/transformations/cms/PrepareCMSReconfiguration.java @@ -101,9 +101,7 @@ public class PrepareCMSReconfiguration .map(prev.directory::peerId) .collect(Collectors.toSet()); - Set withoutReplaced = new HashSet<>(currentCms); - withoutReplaced.remove(toReplace); - Set newCms = placementStrategy.reconfigure(withoutReplaced, prev); + Set newCms = placementStrategy.reconfigure(prev); if (newCms.equals(currentCms)) { logger.info("Proposed CMS reconfiguration resulted in no required modifications at epoch {}", prev.epoch.getEpoch()); @@ -171,7 +169,7 @@ public class PrepareCMSReconfiguration .map(prev.directory::peerId) .collect(Collectors.toSet()); - Set newCms = placementStrategy.reconfigure(currentCms, prev); + Set newCms = placementStrategy.reconfigure(prev); if (newCms.equals(currentCms)) { logger.info("Proposed CMS reconfiguration resulted in no required modifications at epoch {}", prev.epoch.getEpoch()); @@ -235,18 +233,6 @@ public class PrepareCMSReconfiguration return new Diff(additions, removals); } - public static boolean needsReconfiguration(ClusterMetadata metadata) - { - CMSPlacementStrategy placementStrategy = CMSPlacementStrategy.fromReplicationParams(ReplicationParams.meta(metadata), nodeId -> true); - Set currentCms = metadata.fullCMSMembers() - .stream() - .map(metadata.directory::peerId) - .collect(Collectors.toSet()); - - Set newCms = placementStrategy.reconfigure(currentCms, metadata); - return !currentCms.equals(newCms); - } - public static class Diff { public static final Serializer serializer = new Serializer(); diff --git a/src/java/org/apache/cassandra/utils/btree/BTreeMultimap.java b/src/java/org/apache/cassandra/utils/btree/BTreeMultimap.java index aefa94a8dc..3f5a9438e6 100644 --- a/src/java/org/apache/cassandra/utils/btree/BTreeMultimap.java +++ b/src/java/org/apache/cassandra/utils/btree/BTreeMultimap.java @@ -19,6 +19,7 @@ package org.apache.cassandra.utils.btree; import java.util.Collection; +import java.util.Collections; import java.util.Comparator; import java.util.HashSet; import java.util.Map; @@ -138,7 +139,10 @@ public class BTreeMultimap implements Multimap { if (k == null) return null; - return map.get(k); + Collection value = map.get(k); + if (value == null) + return Collections.emptySet(); + return value; } @Override diff --git a/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/AssassinateCMSNodeTest.java b/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/AssassinateCMSNodeTest.java new file mode 100644 index 0000000000..e84b51ff80 --- /dev/null +++ b/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/AssassinateCMSNodeTest.java @@ -0,0 +1,54 @@ +/* + * 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.cassandra.distributed.test.hostreplacement; + +import java.io.IOException; +import java.net.InetSocketAddress; +import java.util.concurrent.ExecutionException; + +import org.junit.Test; + +import org.apache.cassandra.distributed.Cluster; +import org.apache.cassandra.distributed.test.TestBaseImpl; +import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.utils.FBUtilities; + +import static org.apache.cassandra.distributed.api.Feature.GOSSIP; +import static org.apache.cassandra.distributed.api.Feature.NETWORK; +import static org.junit.Assert.assertTrue; + +public class AssassinateCMSNodeTest extends TestBaseImpl +{ + @Test + public void assassinateCMSNodeTest() throws IOException, ExecutionException, InterruptedException + { + try (Cluster cluster = init(builder().withNodes(3) + .withConfig(config -> config.with(GOSSIP).with(NETWORK)) + .start())) + { + cluster.get(1).nodetoolResult("cms", "reconfigure", "3").asserts().success(); + InetSocketAddress toAssassinate = cluster.get(2).broadcastAddress(); + cluster.get(2).shutdown().get(); + cluster.get(1).nodetoolResult("assassinate", toAssassinate.getHostString()).asserts().success(); + cluster.get(1).runOnInstance(() -> assertTrue(ClusterMetadata.current().isCMSMember(FBUtilities.getBroadcastAddressAndPort()))); + cluster.get(3).runOnInstance(() -> assertTrue(ClusterMetadata.current().isCMSMember(FBUtilities.getBroadcastAddressAndPort()))); + cluster.get(1).nodetoolResult("cms").asserts().success(); + } + } +} diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java index d5c25fe6d2..89253b72b8 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/MetadataChangeSimulationTest.java @@ -550,8 +550,7 @@ public class MetadataChangeSimulationTest extends CMSTestBase } } - Set newCms = CMSConfigurationStrategy.reconfigure(sut.service.metadata().directory.toNodeIds(sut.service.metadata().fullCMSMembers()), - sut.service.metadata()); + Set newCms = CMSConfigurationStrategy.reconfigure(sut.service.metadata()); ClusterMetadata metadata = sut.service.metadata(); diff --git a/test/unit/org/apache/cassandra/locator/MetaStrategyTest.java b/test/unit/org/apache/cassandra/locator/MetaStrategyTest.java index a357fc88a4..d90cbd2464 100644 --- a/test/unit/org/apache/cassandra/locator/MetaStrategyTest.java +++ b/test/unit/org/apache/cassandra/locator/MetaStrategyTest.java @@ -113,21 +113,21 @@ public class MetaStrategyTest CMSPlacementStrategy placementStrategy = new CMSPlacementStrategy.DatacenterAware(rf, (cd, n) -> true); Assert.assertEquals(nodeIds(metadata.directory, 1, 2, 4, 5, 7, 8), - placementStrategy.reconfigure(Collections.EMPTY_SET, metadata)); + placementStrategy.reconfigure(metadata)); Assert.assertEquals(nodeIds(metadata.directory, 1, 2, 4, 5, 7, 8), - placementStrategy.reconfigure(nodeIds(metadata.directory, 3, 6, 9), metadata)); + placementStrategy.reconfigure(metadata)); placementStrategy = new CMSPlacementStrategy.DatacenterAware(rf, (cd, n) -> !n.equals(metadata.directory.peerId(addr(2).broadcastAddress)) && !n.equals(metadata.directory.peerId(addr(2).broadcastAddress))); Assert.assertEquals(nodeIds(metadata.directory, 1, 3, 4, 5, 7, 8), - placementStrategy.reconfigure(Collections.EMPTY_SET, metadata)); + placementStrategy.reconfigure(metadata)); Assert.assertEquals(nodeIds(metadata.directory, 1, 3, 4, 5, 7, 8), - placementStrategy.reconfigure(nodeIds(metadata.directory, 3, 6, 9), metadata)); + placementStrategy.reconfigure(metadata)); } public static Set nodeIds(Directory directory, int... addrs) throws UnknownHostException diff --git a/test/unit/org/apache/cassandra/utils/btree/BTreeMultimapTest.java b/test/unit/org/apache/cassandra/utils/btree/BTreeMultimapTest.java index c3b980ac37..519f0d9dd4 100644 --- a/test/unit/org/apache/cassandra/utils/btree/BTreeMultimapTest.java +++ b/test/unit/org/apache/cassandra/utils/btree/BTreeMultimapTest.java @@ -19,6 +19,7 @@ package org.apache.cassandra.utils.btree; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Random; @@ -54,6 +55,7 @@ public class BTreeMultimapTest map = map.without("hello", 125); assertEquals(0, map.size()); assertFalse(map.containsKey("hello")); + assertEquals(Collections.emptySet(), map.get("non-existing")); } @Test