mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-4.1' into trunk
This commit is contained in:
commit
158875858c
|
|
@ -103,6 +103,7 @@ Merged from 4.1:
|
||||||
* Streaming progress virtual table lock contention can trigger TCP_USER_TIMEOUT and fail streaming (CASSANDRA-18110)
|
* Streaming progress virtual table lock contention can trigger TCP_USER_TIMEOUT and fail streaming (CASSANDRA-18110)
|
||||||
* Fix perpetual load of denylist on read in cases where denylist can never be loaded (CASSANDRA-18116)
|
* Fix perpetual load of denylist on read in cases where denylist can never be loaded (CASSANDRA-18116)
|
||||||
Merged from 4.0:
|
Merged from 4.0:
|
||||||
|
* Add safeguard so cleanup fails when node has pending ranges (CASSANDRA-16418)
|
||||||
* Fix legacy clustering serialization for paging with compact storage (CASSANDRA-17507)
|
* Fix legacy clustering serialization for paging with compact storage (CASSANDRA-17507)
|
||||||
* Add support for python 3.11 (CASSANDRA-18088)
|
* Add support for python 3.11 (CASSANDRA-18088)
|
||||||
* Fix formatting of duration in cqlsh (CASSANDRA-18141)
|
* Fix formatting of duration in cqlsh (CASSANDRA-18141)
|
||||||
|
|
@ -157,7 +158,6 @@ Merged from 3.0:
|
||||||
* Revert removal of withBufferSizeInMB(int size) in CQLSSTableWriter.Builder class and deprecate it in favor of withBufferSizeInMiB(int size) (CASSANDRA-17675)
|
* Revert removal of withBufferSizeInMB(int size) in CQLSSTableWriter.Builder class and deprecate it in favor of withBufferSizeInMiB(int size) (CASSANDRA-17675)
|
||||||
* Remove expired snapshots of dropped tables after restart (CASSANDRA-17619)
|
* Remove expired snapshots of dropped tables after restart (CASSANDRA-17619)
|
||||||
Merged from 4.0:
|
Merged from 4.0:
|
||||||
* Fix sstable loading of keyspaces named snapshots or backups (CASSANDRA-14013)
|
|
||||||
* Restore internode custom tracing on 4.0's new messaging system (CASSANDRA-17981)
|
* Restore internode custom tracing on 4.0's new messaging system (CASSANDRA-17981)
|
||||||
* Harden parsing of boolean values in CQL in PropertyDefinitions (CASSANDRA-17878)
|
* Harden parsing of boolean values in CQL in PropertyDefinitions (CASSANDRA-17878)
|
||||||
* Fix possible race condition on repair snapshots (CASSANDRA-17955)
|
* Fix possible race condition on repair snapshots (CASSANDRA-17955)
|
||||||
|
|
|
||||||
|
|
@ -569,11 +569,7 @@ public class CompactionManager implements CompactionManagerMBean
|
||||||
{
|
{
|
||||||
assert !cfStore.isIndex();
|
assert !cfStore.isIndex();
|
||||||
Keyspace keyspace = cfStore.keyspace;
|
Keyspace keyspace = cfStore.keyspace;
|
||||||
if (!StorageService.instance.isJoined())
|
|
||||||
{
|
|
||||||
logger.info("Cleanup cannot run before a node has joined the ring");
|
|
||||||
return AllSSTableOpStatus.ABORTED;
|
|
||||||
}
|
|
||||||
// if local ranges is empty, it means no data should remain
|
// if local ranges is empty, it means no data should remain
|
||||||
final RangesAtEndpoint replicas = StorageService.instance.getLocalReplicas(keyspace.getName());
|
final RangesAtEndpoint replicas = StorageService.instance.getLocalReplicas(keyspace.getName());
|
||||||
final Set<Range<Token>> allRanges = replicas.ranges();
|
final Set<Range<Token>> allRanges = replicas.ranges();
|
||||||
|
|
|
||||||
|
|
@ -168,6 +168,7 @@ import static org.apache.cassandra.service.ActiveRepairService.*;
|
||||||
import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis;
|
import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis;
|
||||||
import static org.apache.cassandra.utils.Clock.Global.nanoTime;
|
import static org.apache.cassandra.utils.Clock.Global.nanoTime;
|
||||||
import static org.apache.cassandra.utils.FBUtilities.now;
|
import static org.apache.cassandra.utils.FBUtilities.now;
|
||||||
|
import static org.apache.cassandra.utils.FBUtilities.getBroadcastAddressAndPort;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* This abstraction contains the token/identifier of this node
|
* This abstraction contains the token/identifier of this node
|
||||||
|
|
@ -3863,6 +3864,9 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
||||||
if (SchemaConstants.isLocalSystemKeyspace(keyspaceName))
|
if (SchemaConstants.isLocalSystemKeyspace(keyspaceName))
|
||||||
throw new RuntimeException("Cleanup of the system keyspace is neither necessary nor wise");
|
throw new RuntimeException("Cleanup of the system keyspace is neither necessary nor wise");
|
||||||
|
|
||||||
|
if (tokenMetadata.getPendingRanges(keyspaceName, getBroadcastAddressAndPort()).size() > 0)
|
||||||
|
throw new RuntimeException("Node is involved in cluster membership changes. Not safe to run cleanup.");
|
||||||
|
|
||||||
CompactionManager.AllSSTableOpStatus status = CompactionManager.AllSSTableOpStatus.SUCCESSFUL;
|
CompactionManager.AllSSTableOpStatus status = CompactionManager.AllSSTableOpStatus.SUCCESSFUL;
|
||||||
for (ColumnFamilyStore cfStore : getValidColumnFamilies(false, false, keyspaceName, tables))
|
for (ColumnFamilyStore cfStore : getValidColumnFamilies(false, false, keyspaceName, tables))
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -73,6 +73,18 @@ public class GossipHelper
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static InstanceAction statusToDecommission(IInvokableInstance newNode)
|
||||||
|
{
|
||||||
|
return (instance) ->
|
||||||
|
{
|
||||||
|
changeGossipState(instance,
|
||||||
|
newNode,
|
||||||
|
Arrays.asList(tokens(newNode),
|
||||||
|
statusLeaving(newNode),
|
||||||
|
statusWithPortLeaving(newNode)));
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
public static InstanceAction statusToNormal(IInvokableInstance peer)
|
public static InstanceAction statusToNormal(IInvokableInstance peer)
|
||||||
{
|
{
|
||||||
return (target) ->
|
return (target) ->
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,111 @@
|
||||||
|
/*
|
||||||
|
* 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.ring;
|
||||||
|
|
||||||
|
import org.junit.Test;
|
||||||
|
|
||||||
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
|
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
||||||
|
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
||||||
|
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||||
|
import org.apache.cassandra.distributed.api.NodeToolResult;
|
||||||
|
import org.apache.cassandra.distributed.shared.NetworkTopology;
|
||||||
|
import org.apache.cassandra.distributed.test.TestBaseImpl;
|
||||||
|
|
||||||
|
import static org.junit.Assert.assertEquals;
|
||||||
|
import static org.apache.cassandra.distributed.action.GossipHelper.statusToBootstrap;
|
||||||
|
import static org.apache.cassandra.distributed.action.GossipHelper.statusToDecommission;
|
||||||
|
import static org.apache.cassandra.distributed.action.GossipHelper.withProperty;
|
||||||
|
import static org.apache.cassandra.distributed.test.ring.BootstrapTest.populate;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
||||||
|
import static org.apache.cassandra.distributed.api.TokenSupplier.evenlyDistributedTokens;
|
||||||
|
|
||||||
|
public class CleanupFailureTest extends TestBaseImpl
|
||||||
|
{
|
||||||
|
@Test
|
||||||
|
public void cleanupDuringDecommissionTest() throws Throwable
|
||||||
|
{
|
||||||
|
try (Cluster cluster = builder().withNodes(2)
|
||||||
|
.withTokenSupplier(evenlyDistributedTokens(2))
|
||||||
|
.withNodeIdTopology(NetworkTopology.singleDcNetworkTopology(2, "dc0", "rack0"))
|
||||||
|
.withConfig(config -> config.with(NETWORK, GOSSIP))
|
||||||
|
.start())
|
||||||
|
{
|
||||||
|
IInvokableInstance nodeToDecommission = cluster.get(1);
|
||||||
|
IInvokableInstance nodeToRemainInCluster = cluster.get(2);
|
||||||
|
|
||||||
|
// Start decomission on nodeToDecommission
|
||||||
|
cluster.forEach(statusToDecommission(nodeToDecommission));
|
||||||
|
|
||||||
|
// Add data to cluster while node is decomissioning
|
||||||
|
int NUM_ROWS = 100;
|
||||||
|
populate(cluster, 0, NUM_ROWS, 1, 1, ConsistencyLevel.ONE);
|
||||||
|
cluster.forEach(c -> c.flush(KEYSPACE));
|
||||||
|
|
||||||
|
// Check data before cleanup on nodeToRemainInCluster
|
||||||
|
assertEquals(100, nodeToRemainInCluster.executeInternal("SELECT * FROM " + KEYSPACE + ".tbl").length);
|
||||||
|
|
||||||
|
// Run cleanup on nodeToRemainInCluster
|
||||||
|
NodeToolResult result = nodeToRemainInCluster.nodetoolResult("cleanup");
|
||||||
|
result.asserts().failure();
|
||||||
|
result.asserts().stderrContains("Node is involved in cluster membership changes. Not safe to run cleanup.");
|
||||||
|
|
||||||
|
// Check data after cleanup on nodeToRemainInCluster
|
||||||
|
assertEquals(100, nodeToRemainInCluster.executeInternal("SELECT * FROM " + KEYSPACE + ".tbl").length);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void cleanupDuringBootstrapTest() throws Throwable
|
||||||
|
{
|
||||||
|
int originalNodeCount = 1;
|
||||||
|
int expandedNodeCount = originalNodeCount + 1;
|
||||||
|
|
||||||
|
try (Cluster cluster = builder().withNodes(originalNodeCount)
|
||||||
|
.withTokenSupplier(evenlyDistributedTokens(expandedNodeCount))
|
||||||
|
.withNodeIdTopology(NetworkTopology.singleDcNetworkTopology(expandedNodeCount, "dc0", "rack0"))
|
||||||
|
.withConfig(config -> config.with(NETWORK, GOSSIP))
|
||||||
|
.start())
|
||||||
|
{
|
||||||
|
IInstanceConfig config = cluster.newInstanceConfig();
|
||||||
|
IInvokableInstance bootstrappingNode = cluster.bootstrap(config);
|
||||||
|
withProperty("cassandra.join_ring", false,
|
||||||
|
() -> bootstrappingNode.startup(cluster));
|
||||||
|
|
||||||
|
// Start decomission on bootstrappingNode
|
||||||
|
cluster.forEach(statusToBootstrap(bootstrappingNode));
|
||||||
|
|
||||||
|
// Add data to cluster while node is bootstrapping
|
||||||
|
int NUM_ROWS = 100;
|
||||||
|
populate(cluster, 0, NUM_ROWS, 1, 2, ConsistencyLevel.ONE);
|
||||||
|
cluster.forEach(c -> c.flush(KEYSPACE));
|
||||||
|
|
||||||
|
// Check data before cleanup on bootstrappingNode
|
||||||
|
assertEquals(NUM_ROWS, bootstrappingNode.executeInternal("SELECT * FROM " + KEYSPACE + ".tbl").length);
|
||||||
|
|
||||||
|
// Run cleanup on bootstrappingNode
|
||||||
|
NodeToolResult result = bootstrappingNode.nodetoolResult("cleanup");
|
||||||
|
result.asserts().stderrContains("Node is involved in cluster membership changes. Not safe to run cleanup.");
|
||||||
|
|
||||||
|
// Check data after cleanup on bootstrappingNode
|
||||||
|
assertEquals(NUM_ROWS, bootstrappingNode.executeInternal("SELECT * FROM " + KEYSPACE + ".tbl").length);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -117,9 +117,11 @@ public class CleanupTest
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void testCleanup() throws ExecutionException, InterruptedException
|
public void testCleanup() throws ExecutionException, InterruptedException, UnknownHostException
|
||||||
{
|
{
|
||||||
StorageService.instance.getTokenMetadata().clearUnsafe();
|
TokenMetadata tmd = StorageService.instance.getTokenMetadata();
|
||||||
|
tmd.clearUnsafe();
|
||||||
|
tmd.updateNormalToken(token(new byte[]{ 50 }), InetAddressAndPort.getByName("127.0.0.1"));
|
||||||
|
|
||||||
Keyspace keyspace = Keyspace.open(KEYSPACE1);
|
Keyspace keyspace = Keyspace.open(KEYSPACE1);
|
||||||
ColumnFamilyStore cfs = keyspace.getColumnFamilyStore(CF_STANDARD1);
|
ColumnFamilyStore cfs = keyspace.getColumnFamilyStore(CF_STANDARD1);
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue