Merge branch 'cassandra-3.0' into cassandra-3.11

This commit is contained in:
Ekaterina Dimitrova 2021-03-10 20:50:11 -05:00
commit b654000936
21 changed files with 1029 additions and 98 deletions

View File

@ -2,7 +2,7 @@ version: 2
jobs:
j8_jvm_upgrade_dtests:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -94,7 +94,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
build:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -175,7 +175,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_dtests-no-vnodes:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -233,7 +233,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_upgradetests-no-vnodes:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -332,7 +332,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
utests_stress:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -378,7 +378,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_unit_tests:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -470,7 +470,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_dtests-with-vnodes:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -528,7 +528,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_jvm_dtests:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -620,7 +620,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
utests_long:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -666,7 +666,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
utests_compression:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l
@ -758,7 +758,7 @@ jobs:
- JDK_HOME: /usr/lib/jvm/java-8-openjdk-amd64
j8_dtest_jars_build:
docker:
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210105
- image: apache/cassandra-testing-ubuntu2004-java11-w-dependencies:20210304
resource_class: medium
working_directory: ~/
shell: /bin/bash -eo pipefail -l

View File

@ -1,6 +1,7 @@
3.11.11
* Upgrade jackson-databind to 2.9.10.8 (CASSANDRA-16462)
Merged from 3.0:
* Refuse DROP COMPACT STORAGE if some 2.x sstables are in use (CASSANDRA-15897)
* Fix ColumnFilter::toString not returning a valid CQL fragment (CASSANDRA-16483)
* Fix ColumnFilter behaviour to prevent digest mitmatches during upgrades (CASSANDRA-16415)
* Avoid pushing schema mutations when setting up distributed system keyspaces locally (CASSANDRA-16387)

View File

@ -426,6 +426,7 @@
<dependency groupId="org.mockito" artifactId="mockito-core" version="3.2.4" />
<dependency groupId="org.apache.cassandra" artifactId="dtest-api" version="0.0.7" />
<dependency groupId="org.reflections" artifactId="reflections" version="0.9.12" />
<dependency groupId="org.quicktheories" artifactId="quicktheories" version="0.25" />
<dependency groupId="org.apache.rat" artifactId="apache-rat" version="0.10">
<exclusion groupId="commons-lang" artifactId="commons-lang"/>
</dependency>
@ -486,6 +487,8 @@
<dependency groupId="com.github.ben-manes.caffeine" artifactId="caffeine" version="2.2.6" />
<dependency groupId="org.jctools" artifactId="jctools-core" version="1.2.1"/>
<dependency groupId="org.ow2.asm" artifactId="asm" version="5.0.4" />
<!-- when updating assertj, make sure to also update the corresponding junit-bom dependency -->
<dependency groupId="org.assertj" artifactId="assertj-core" version="3.15.0"/>
</dependencyManagement>
<developer id="adelapena" name="Andres de la Peña"/>
<developer id="alakshman" name="Avinash Lakshman"/>
@ -543,6 +546,7 @@
version="${version}"/>
<dependency groupId="junit" artifactId="junit"/>
<dependency groupId="org.mockito" artifactId="mockito-core" />
<dependency groupId="org.quicktheories" artifactId="quicktheories" />
<dependency groupId="org.apache.cassandra" artifactId="dtest-api" />
<dependency groupId="org.reflections" artifactId="reflections" />
<dependency groupId="org.apache.rat" artifactId="apache-rat"/>
@ -563,6 +567,10 @@
<dependency groupId="org.openjdk.jmh" artifactId="jmh-generator-annprocess"/>
<dependency groupId="net.ju-n.compile-command-annotations" artifactId="compile-command-annotations"/>
<dependency groupId="org.apache.ant" artifactId="ant-junit" version="1.9.4" />
<!-- adding this dependency is necessary for assertj. When updating assertj, need to also update the version of
this that the new assertj's `assertj-parent-pom` depends on. -->
<dependency groupId="org.junit" artifactId="junit-bom" version="5.6.0" type="pom"/>
<dependency groupId="org.assertj" artifactId="assertj-core"/>
</artifact:pom>
<!-- this build-deps-pom-sources "artifact" is the same as build-deps-pom but only with those
artifacts that have "-source.jar" files -->
@ -588,6 +596,7 @@
<dependency groupId="org.openjdk.jmh" artifactId="jmh-generator-annprocess"/>
<dependency groupId="net.ju-n.compile-command-annotations" artifactId="compile-command-annotations"/>
<dependency groupId="org.apache.ant" artifactId="ant-junit" version="1.9.4" />
<dependency groupId="org.assertj" artifactId="assertj-core"/>
</artifact:pom>
<artifact:pom id="coverage-deps-pom"
@ -1329,6 +1338,8 @@
<jvmarg value="-Dcassandra.testtag=@{testtag}"/>
<jvmarg value="-Dcassandra.keepBriefBrief=${cassandra.keepBriefBrief}" />
<jvmarg value="-Dcassandra.strict.runtime.checks=true" />
<!-- disable shrinks in quicktheories CASSANDRA-15554 -->
<jvmarg value="-DQT_SHRINKS=0"/>
<optjvmargs/>
<classpath>
<pathelement path="${java.class.path}"/>

View File

@ -17,11 +17,17 @@
*/
package org.apache.cassandra.cql3.statements;
import java.net.InetAddress;
import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import com.google.common.base.Splitter;
import com.google.common.collect.Iterables;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.auth.Permission;
import org.apache.cassandra.config.*;
import org.apache.cassandra.cql3.*;
@ -32,19 +38,28 @@ import org.apache.cassandra.db.marshal.BytesType;
import org.apache.cassandra.db.marshal.EmptyType;
import org.apache.cassandra.db.view.View;
import org.apache.cassandra.exceptions.*;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.gms.ApplicationState;
import org.apache.cassandra.gms.Gossiper;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.schema.IndexMetadata;
import org.apache.cassandra.schema.Indexes;
import org.apache.cassandra.schema.TableParams;
import org.apache.cassandra.service.ClientState;
import org.apache.cassandra.service.MigrationManager;
import org.apache.cassandra.service.QueryState;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.transport.Event;
import org.apache.cassandra.utils.NoSpamLogger;
import static java.lang.String.format;
import static org.apache.cassandra.thrift.ThriftValidation.validateColumnFamily;
public class AlterTableStatement extends SchemaAlteringStatement
{
private static final Logger logger = LoggerFactory.getLogger(AlterTableStatement.class);
private static final NoSpamLogger noSpamLogger = NoSpamLogger.getLogger(logger, 5L, TimeUnit.MINUTES);
public enum Type
{
ADD, ALTER, DROP, DROP_COMPACT_STORAGE, OPTS, RENAME
@ -241,24 +256,24 @@ public class AlterTableStatement extends SchemaAlteringStatement
switch (def.kind)
{
case PARTITION_KEY:
case CLUSTERING:
throw new InvalidRequestException(String.format("Cannot drop PRIMARY KEY part %s", columnName));
case REGULAR:
case STATIC:
ColumnDefinition toDelete = null;
for (ColumnDefinition columnDef : cfm.partitionColumns())
{
if (columnDef.name.equals(columnName))
{
toDelete = columnDef;
break;
}
}
assert toDelete != null;
cfm.removeColumnDefinition(toDelete);
cfm.recordColumnDrop(toDelete, deleteTimestamp == null ? queryState.getTimestamp() : deleteTimestamp);
break;
case PARTITION_KEY:
case CLUSTERING:
throw new InvalidRequestException(String.format("Cannot drop PRIMARY KEY part %s", columnName));
case REGULAR:
case STATIC:
ColumnDefinition toDelete = null;
for (ColumnDefinition columnDef : cfm.partitionColumns())
{
if (columnDef.name.equals(columnName))
{
toDelete = columnDef;
break;
}
}
assert toDelete != null;
cfm.removeColumnDefinition(toDelete);
cfm.recordColumnDrop(toDelete, deleteTimestamp == null ? queryState.getTimestamp() : deleteTimestamp);
break;
}
// If the dropped column is required by any secondary indexes
@ -278,22 +293,17 @@ public class AlterTableStatement extends SchemaAlteringStatement
}
if (!Iterables.isEmpty(views))
throw new InvalidRequestException(String.format("Cannot drop column %s on base table %s with materialized views.",
throw new InvalidRequestException(String.format("Cannot drop column %s on base table %s with materialized views.",
columnName.toString(),
columnFamily()));
}
break;
case DROP_COMPACT_STORAGE:
if (!meta.isCompactTable())
throw new InvalidRequestException("Cannot DROP COMPACT STORAGE on table without COMPACT STORAGE");
// TODO: Global check of the sstables to be added as part of CASSANDRA-15897.
// Currently this is only a local check of the SSTables versions
for (SSTableReader ssTableReader : Keyspace.open(keyspace()).getColumnFamilyStore(columnFamily()).getLiveSSTables())
{
if (!ssTableReader.descriptor.version.isLatestVersion())
throw new InvalidRequestException("Cannot DROP COMPACT STORAGE until all SSTables are upgraded, please run `nodetool upgradesstables` first.");
}
validateCanDropCompactStorage();
cfm = meta.asNonCompact();
break;
@ -354,6 +364,75 @@ public class AlterTableStatement extends SchemaAlteringStatement
return new Event.SchemaChange(Event.SchemaChange.Change.UPDATED, Event.SchemaChange.Target.TABLE, keyspace(), columnFamily());
}
/**
* Throws if DROP COMPACT STORAGE cannot be used (yet) because the cluster is not sufficiently upgraded. To be able
* to use DROP COMPACT STORAGE, we need to ensure that no pre-3.0 sstables exists in the cluster, as we won't be
* able to read them anymore once COMPACT STORAGE is dropped (see CASSANDRA-15897). In practice, this method checks
* 3 things:
* 1) that all nodes are on 3.0+. We need this because 2.x nodes don't advertise their sstable versions.
* 2) for 3.0+, we use the new (CASSANDRA-15897) sstables versions set gossiped by all nodes to ensure all
* sstables have been upgraded cluster-wise.
* 3) if the cluster still has some 3.0 nodes that predate CASSANDRA-15897, we will not have the sstable versions
* for them. In that case, we also refuse DROP COMPACT (even though it may well be safe at this point) and ask
* the user to upgrade all nodes.
*/
private void validateCanDropCompactStorage()
{
Set<InetAddress> before3 = new HashSet<>();
Set<InetAddress> preC15897nodes = new HashSet<>();
Set<InetAddress> with2xSStables = new HashSet<>();
Splitter onComma = Splitter.on(',').omitEmptyStrings().trimResults();
for (InetAddress node : StorageService.instance.getTokenMetadata().getAllEndpoints())
{
if (MessagingService.instance().getVersion(node) < MessagingService.VERSION_30)
{
before3.add(node);
continue;
}
String sstableVersionsString = Gossiper.instance.getApplicationState(node, ApplicationState.SSTABLE_VERSIONS);
if (sstableVersionsString == null)
{
preC15897nodes.add(node);
continue;
}
try
{
boolean has2xSStables = onComma.splitToList(sstableVersionsString)
.stream()
.map(VersionAndType::fromString)
.anyMatch(v -> !v.version().storeRows());
if (has2xSStables)
with2xSStables.add(node);
}
catch (IllegalArgumentException e)
{
// Means VersionType::fromString didn't parse a version correctly. Which shouldn't happen, we shouldn't
// have garbage in Gossip. But crashing the request is not ideal, so we log the error but ignore the
// node otherwise.
noSpamLogger.error("Unexpected error parsing sstable versions from gossip for {} (gossiped value " +
"is '{}'). This is a bug and should be reported. Cannot ensure that {} has no " +
"non-upgraded 2.x sstables anymore. If after this DROP COMPACT STORAGE some old " +
"sstables cannot be read anymore, please use `upgradesstables` with the " +
"`--force-compact-storage-on` option.", node, sstableVersionsString, node);
}
}
if (!before3.isEmpty())
throw new InvalidRequestException(format("Cannot DROP COMPACT STORAGE as some nodes in the cluster (%s) " +
"are not on 3.0+ yet. Please upgrade those nodes and run " +
"`upgradesstables` before retrying.", before3));
if (!preC15897nodes.isEmpty())
throw new InvalidRequestException(format("Cannot guarantee that DROP COMPACT STORAGE is safe as some nodes " +
"in the cluster (%s) do not have https://issues.apache.org/jira/browse/CASSANDRA-15897. " +
"Please upgrade those nodes and retry.", preC15897nodes));
if (!with2xSStables.isEmpty())
throw new InvalidRequestException(format("Cannot DROP COMPACT STORAGE as some nodes in the cluster (%s) " +
"has some non-upgraded 2.x sstables. Please run `upgradesstables` " +
"on those nodes before retrying", with2xSStables));
}
public String toString()
{
return String.format("AlterTableStatement(name=%s, type=%s)",

View File

@ -426,6 +426,10 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
initialMemtable = new Memtable(new AtomicReference<>(CommitLog.instance.getCurrentPosition()), this);
data = new Tracker(initialMemtable, loadSSTables);
// Note that this needs to happen before we load the first sstables, or the global sstable tracker will not
// be notified on the initial loading.
data.subscribe(StorageService.instance.sstablesTracker);
Collection<SSTableReader> sstables = null;
// scan for sstables corresponding to this cf and load them
if (data.loadsstables)

View File

@ -582,7 +582,8 @@ public class CompactionStrategyManager implements INotificationConsumer
{
if (notification instanceof SSTableAddedNotification)
{
handleFlushNotification(((SSTableAddedNotification) notification).added);
SSTableAddedNotification flushedNotification = (SSTableAddedNotification) notification;
handleFlushNotification(flushedNotification.added);
}
else if (notification instanceof SSTableListChangedNotification)
{

View File

@ -185,17 +185,12 @@ public class Tracker
public void addInitialSSTables(Iterable<SSTableReader> sstables)
{
addInitialSSTablesWithoutUpdatingSize(sstables);
maybeFail(updateSizeTracking(emptySet(), sstables, null));
// no notifications or backup necessary
addSSTablesInternal(sstables, true, false, true);
}
public void addInitialSSTablesWithoutUpdatingSize(Iterable<SSTableReader> sstables)
{
if (!isDummy())
setupOnline(sstables);
apply(updateLiveSet(emptySet(), sstables));
// no notifications or backup necessary
addSSTablesInternal(sstables, true, false, false);
}
public void updateInitialSSTableSize(Iterable<SSTableReader> sstables)
@ -205,9 +200,22 @@ public class Tracker
public void addSSTables(Iterable<SSTableReader> sstables)
{
addInitialSSTables(sstables);
maybeIncrementallyBackup(sstables);
notifyAdded(sstables);
addSSTablesInternal(sstables, false, true, true);
}
private void addSSTablesInternal(Iterable<SSTableReader> sstables,
boolean isInitialSSTables,
boolean maybeIncrementallyBackup,
boolean updateSize)
{
if (!isDummy())
setupOnline(sstables);
apply(updateLiveSet(emptySet(), sstables));
if(updateSize)
maybeFail(updateSizeTracking(emptySet(), sstables, null));
if (maybeIncrementallyBackup)
maybeIncrementallyBackup(sstables);
notifyAdded(sstables, isInitialSSTables);
}
/** (Re)initializes the tracker, purging all references. */
@ -370,7 +378,7 @@ public class Tracker
notifyDiscarded(memtable);
// TODO: if we're invalidated, should we notifyadded AND removed, or just skip both?
fail = notifyAdded(sstables, fail);
fail = notifyAdded(sstables, false, fail);
if (!isDummy() && !cfstore.isValid())
dropSSTables();
@ -428,9 +436,14 @@ public class Tracker
return accumulate;
}
Throwable notifyAdded(Iterable<SSTableReader> added, Throwable accumulate)
Throwable notifyAdded(Iterable<SSTableReader> added, boolean isInitialSSTables, Throwable accumulate)
{
INotification notification = new SSTableAddedNotification(added);
INotification notification;
if (!isInitialSSTables)
notification = new SSTableAddedNotification(added);
else
notification = new InitialSSTableAddedNotification(added);
for (INotificationConsumer subscriber : subscribers)
{
try
@ -445,9 +458,9 @@ public class Tracker
return accumulate;
}
public void notifyAdded(Iterable<SSTableReader> added)
void notifyAdded(Iterable<SSTableReader> added, boolean isInitialSSTables)
{
maybeFail(notifyAdded(added, null));
maybeFail(notifyAdded(added, isInitialSSTables, null));
}
public void notifySSTableRepairedStatusChanged(Collection<SSTableReader> repairStatusesChanged)

View File

@ -17,6 +17,14 @@
*/
package org.apache.cassandra.gms;
/**
* The various "states" exchanged through Gossip.
*
* <p><b>Important Note:</b> Gossip uses the ordinal of this enum in the messages it exchanges, so values in that enum
* should <i>not</i> be re-ordered or removed. The end of this enum should also always include some "padding" so that
* if newer versions add new states, old nodes that don't know about those new states don't "break" deserializing those
* states.
*/
public enum ApplicationState
{
STATUS,
@ -34,11 +42,22 @@ public enum ApplicationState
HOST_ID,
TOKENS,
RPC_READY,
// pad to allow adding new states to existing cluster
// We added SSTABLE_VERSIONS in CASSANDRA-15897 in 3.0, and at the time, 3 more ApplicationState had been added
// to newer versions, so we skipped the first 3 of our original padding to ensure SSTABLE_VERSIONS can preserve
// its ordinal accross versions.
X1,
X2,
X3,
X4,
/**
* The set of sstable versions on this node. This will usually be only the "current" sstable format (the one with
* which new sstables are written), but may contain more on newly upgraded nodes before `upgradesstable` has been
* run.
*
* <p>The value (a set of sstable {@link org.apache.cassandra.io.sstable.format.VersionAndType}) is serialized as
* a comma-separated list.
**/
SSTABLE_VERSIONS,
// pad to allow adding new states to existing cluster
X5,
X6,
X7,

View File

@ -953,6 +953,24 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean
return UUID.fromString(epStates.get(endpoint).getApplicationState(ApplicationState.HOST_ID).value);
}
/**
* The value for the provided application state for the provided endpoint as currently known by this Gossip instance.
*
* @param endpoint the endpoint from which to get the endpoint state.
* @param state the endpoint state to get.
* @return the value of the application state {@code state} for {@code endpoint}, or {@code null} if either
* {@code endpoint} is not known by Gossip or has no value for {@code state}.
*/
public String getApplicationState(InetAddress endpoint, ApplicationState state)
{
EndpointState epState = endpointStateMap.get(endpoint);
if (epState == null)
return null;
VersionedValue value = epState.getApplicationState(state);
return value == null ? null : value.value;
}
EndpointState getStateForVersionBiggerThan(InetAddress forEndpoint, int version)
{
EndpointState epState = endpointStateMap.get(forEndpoint);

View File

@ -20,7 +20,9 @@ package org.apache.cassandra.gms;
import java.io.*;
import java.net.InetAddress;
import java.util.Collection;
import java.util.Set;
import java.util.UUID;
import java.util.stream.Collectors;
import static java.nio.charset.StandardCharsets.ISO_8859_1;
@ -31,6 +33,8 @@ import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.net.MessagingService;
@ -274,6 +278,13 @@ public class VersionedValue implements Comparable<VersionedValue>
{
return new VersionedValue(String.valueOf(value));
}
public VersionedValue sstableVersions(Set<VersionAndType> versions)
{
return new VersionedValue(versions.stream()
.map(VersionAndType::toString)
.collect(Collectors.joining(",")));
}
}
private static class VersionedValueSerializer implements IVersionedSerializer<VersionedValue>

View File

@ -358,6 +358,7 @@ public class Descriptor
&& that.generation == this.generation
&& that.ksname.equals(this.ksname)
&& that.cfname.equals(this.cfname)
&& that.version.equals(this.version)
&& that.formatType == this.formatType;
}

View File

@ -0,0 +1,93 @@
/*
* 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.io.sstable.format;
import java.util.List;
import java.util.Objects;
import com.google.common.base.Splitter;
import org.apache.cassandra.io.sstable.Descriptor;
/**
* Groups a sstable {@link Version} with a {@link SSTableFormat.Type}.
*
* <p>Note that both information are currently necessary to identify the exact "format" of an sstable (without having
* its {@link Descriptor}). In particular, while {@link Version} contains its {{@link SSTableFormat}}, you cannot get
* the {{@link SSTableFormat.Type}} from that.
*/
public final class VersionAndType
{
private static final Splitter splitOnDash = Splitter.on('-').omitEmptyStrings().trimResults();
private final Version version;
private final SSTableFormat.Type formatType;
public VersionAndType(Version version, SSTableFormat.Type formatType)
{
this.version = version;
this.formatType = formatType;
}
public Version version()
{
return version;
}
public SSTableFormat.Type formatType()
{
return formatType;
}
@Override
public boolean equals(Object o)
{
if (this == o)
return true;
if (o == null || getClass() != o.getClass())
return false;
VersionAndType that = (VersionAndType) o;
return Objects.equals(version, that.version) &&
formatType == that.formatType;
}
@Override
public int hashCode()
{
return Objects.hash(version, formatType);
}
public static VersionAndType fromString(String versionAndType)
{
List<String> components = splitOnDash.splitToList(versionAndType);
if (components.size() != 2)
throw new IllegalArgumentException("Invalid VersionAndType string: " + versionAndType + " (should be of the form 'big-bc')");
SSTableFormat.Type formatType = SSTableFormat.Type.validate(components.get(0));
Version version = formatType.info.getVersion(components.get(1));
return new VersionAndType(version, formatType);
}
@Override
public String toString()
{
return formatType.name + '-' + version;
}
}

View File

@ -0,0 +1,36 @@
/*
* 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.notifications;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.io.sstable.format.SSTableReader;
/**
* Fired when we load SSTables during the Cassandra initialization phase.
*/
public class InitialSSTableAddedNotification implements INotification
{
/** {@code true} if the addition corresponds to the {@link ColumnFamilyStore} initialization, then the sstables
* are those loaded at startup. */
public final Iterable<SSTableReader> added;
public InitialSSTableAddedNotification(Iterable<SSTableReader> added)
{
this.added = added;
}
}

View File

@ -0,0 +1,284 @@
/*
* 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.service;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArraySet;
import java.util.concurrent.TimeUnit;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Iterables;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.lifecycle.Tracker;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.apache.cassandra.notifications.INotification;
import org.apache.cassandra.notifications.INotificationConsumer;
import org.apache.cassandra.notifications.InitialSSTableAddedNotification;
import org.apache.cassandra.notifications.SSTableAddedNotification;
import org.apache.cassandra.notifications.SSTableDeletingNotification;
import org.apache.cassandra.notifications.SSTableListChangedNotification;
import org.apache.cassandra.utils.NoSpamLogger;
/**
* Tracks all sstables in use on the local node.
*
* <p>Each table tracks its own SSTables in {@link ColumnFamilyStore} (through {@link Tracker}) for most purposes, but
* this class groups information we need on all the sstables the node has.
*/
public class SSTablesGlobalTracker implements INotificationConsumer
{
private static final Logger logger = LoggerFactory.getLogger(SSTablesGlobalTracker.class);
private static final NoSpamLogger noSpamLogger = NoSpamLogger.getLogger(logger, 5L, TimeUnit.MINUTES);
/*
* As of CASSANDRA-15897, the only thing we track here is the set of sstable versions in use.
*
* That set is maintained in `versionsInUse`, an immutable set replaced when changed (so that it can be read safely
* and cheaply). To track when that set changes and needs to be re-computed, we essentially maintain a map of
* sstables versions in use to the number of sstables for that version. But as we know that sstables will
* overwhelmingly be on the "current" version, we special case said current version (makes things cheaper without
* too much complexity here). Those are `sstablesForCurrentVersion` and `sstablesForOtherVersions`.
*
* This would be sufficient if we could guarantee that for any sstable, we only ever have 1 addition notification
* and 1 corresponding remove notification. But sstables can have complex lifecycles and relying on this property
* could be fragile. As a matter of fact, at the time of this writing, the removal notification is sometimes fired
* twice for the same sstable. To keep this component more resilient, we also maintain the set of all sstables for
* which we've received an addition, which allows us to ignore removals for sstables we don't know.
*
* Concurrency handling: the 'allSSTables' set handles concurrency directly as it is updated in all cases. The rest
* of the data structures of this class are only updated together within a synchronized block when handling new
* sstables additions/removals.
*/
private final Set<Descriptor> allSSTables = ConcurrentHashMap.newKeySet();
private final VersionAndType currentVersion;
private int sstablesForCurrentVersion;
private final Map<VersionAndType, Integer> sstablesForOtherVersions = new HashMap<>();
private volatile ImmutableSet<VersionAndType> versionsInUse = ImmutableSet.of();
private final Set<INotificationConsumer> subscribers = new CopyOnWriteArraySet<>();
public SSTablesGlobalTracker(SSTableFormat.Type currentSSTableFormat)
{
this.currentVersion = new VersionAndType(currentSSTableFormat.info.getLatestVersion(), currentSSTableFormat);
}
/**
* The set of all sstable versions currently in use on this node.
*/
public Set<VersionAndType> versionsInUse()
{
return versionsInUse;
}
/**
* Register a new subscriber to this tracker.
*
* Registered subscribers are currently notified when the set of sstable versions in use changes, using a
* {@link SSTablesVersionsInUseChangeNotification}.
*
* @param subscriber the new subscriber to register. If this subscriber is already registered, this method does
* nothing (meaning that even if a subscriber is registered multiple times, it will only be notified once on every
* change).
* @return whether the subscriber was register (so whether it was not already registered).
*/
public boolean register(INotificationConsumer subscriber)
{
return subscribers.add(subscriber);
}
/**
* Unregister a subscriber from this tracker.
*
* @param subscriber the subscriber to unregister. If this subscriber is not registered, this method does nothing.
* @return whether the subscriber was unregistered (so whether it was registered subscriber of this tracker).
*/
public boolean unregister(INotificationConsumer subscriber)
{
return subscribers.remove(subscriber);
}
@Override
public void handleNotification(INotification notification, Object sender)
{
Iterable<Descriptor> removed = removedSSTables(notification);
Iterable<Descriptor> added = addedSSTables(notification);
if (Iterables.isEmpty(removed) && Iterables.isEmpty(added))
return;
boolean triggerUpdate = handleSSTablesChange(removed, added);
if (triggerUpdate)
{
SSTablesVersionsInUseChangeNotification changeNotification = new SSTablesVersionsInUseChangeNotification(versionsInUse);
subscribers.forEach(s -> s.handleNotification(changeNotification, this));
}
}
@VisibleForTesting
boolean handleSSTablesChange(Iterable<Descriptor> removed, Iterable<Descriptor> added)
{
/*
We collect changes to 'sstablesForCurrentVersion' and 'sstablesForOtherVersions' as delta first, and then
apply those delta within a synchronized block below. The goal being to reduce the work done in that
synchronized block.
*/
int currentDelta = 0;
Map<VersionAndType, Integer> othersDelta = null;
/*
Note: we deal with removes first as if a notification both removes and adds, it's a compaction and while
it should never remove and add the same descriptor in practice, doing the remove first is more logical.
*/
for (Descriptor desc : removed)
{
if (!allSSTables.remove(desc))
continue;
VersionAndType version = version(desc);
if (currentVersion.equals(version))
--currentDelta;
else
othersDelta = update(othersDelta, version, -1);
}
for (Descriptor desc : added)
{
if (!allSSTables.add(desc))
continue;
VersionAndType version = version(desc);
if (currentVersion.equals(version))
++currentDelta;
else
othersDelta = update(othersDelta, version, +1);
}
if (currentDelta == 0 && (othersDelta == null))
return false;
/*
Set to true if the set of versions in use is changed by this update. That is, if a version having no
version prior now has some, or if the count for some version reaches 0.
*/
boolean triggerUpdate;
synchronized (this)
{
triggerUpdate = (currentDelta > 0 && sstablesForCurrentVersion == 0)
|| (currentDelta < 0 && sstablesForCurrentVersion <= -currentDelta);
sstablesForCurrentVersion += currentDelta;
sstablesForCurrentVersion = sanitizeSSTablesCount(sstablesForCurrentVersion, currentVersion);
if (othersDelta != null)
{
for (Map.Entry<VersionAndType, Integer> entry : othersDelta.entrySet())
{
VersionAndType version = entry.getKey();
int delta = entry.getValue();
/*
Updates the count, removing the version if it reaches 0 (note: we could use Map#compute for this,
but we wouldn't be able to modify `triggerUpdate` without making it an Object, so we don't bother).
*/
Integer oldValue = sstablesForOtherVersions.get(version);
int newValue = oldValue == null ? delta : oldValue + delta;
newValue = sanitizeSSTablesCount(newValue, version);
triggerUpdate |= oldValue == null || newValue == 0;
if (newValue == 0)
sstablesForOtherVersions.remove(version);
else
sstablesForOtherVersions.put(version, newValue);
}
}
if (triggerUpdate)
versionsInUse = computeVersionsInUse(sstablesForCurrentVersion, currentVersion, sstablesForOtherVersions);
}
return triggerUpdate;
}
private static ImmutableSet<VersionAndType> computeVersionsInUse(int sstablesForCurrentVersion, VersionAndType currentVersion, Map<VersionAndType, Integer> sstablesForOtherVersions)
{
ImmutableSet.Builder<VersionAndType> builder = ImmutableSet.builder();
if (sstablesForCurrentVersion > 0)
builder.add(currentVersion);
builder.addAll(sstablesForOtherVersions.keySet());
return builder.build();
}
private static int sanitizeSSTablesCount(int sstableCount, VersionAndType version)
{
if (sstableCount >= 0)
return sstableCount;
/*
This shouldn't happen and indicate a bug either in the tracking of this class, or on the passed notification.
That said, it's not worth bringing the node down, so we log the problem but otherwise "correct" it.
*/
noSpamLogger.error("Invalid state while handling sstables change notification: the number of sstables for " +
"version {} was computed to {}. This indicate a bug and please report it, but it should " +
"not have adverse consequences.", version, sstableCount, new RuntimeException());
return 0;
}
private static Iterable<Descriptor> addedSSTables(INotification notification)
{
if (notification instanceof SSTableAddedNotification)
return Iterables.transform(((SSTableAddedNotification)notification).added, s -> s.descriptor);
if (notification instanceof SSTableListChangedNotification)
return Iterables.transform(((SSTableListChangedNotification)notification).added, s -> s.descriptor);
if (notification instanceof InitialSSTableAddedNotification)
return Iterables.transform(((InitialSSTableAddedNotification)notification).added, s -> s.descriptor);
else
return Collections.emptyList();
}
private static Iterable<Descriptor> removedSSTables(INotification notification)
{
if (notification instanceof SSTableDeletingNotification)
return Collections.singletonList(((SSTableDeletingNotification)notification).deleting.descriptor);
if (notification instanceof SSTableListChangedNotification)
return Iterables.transform(((SSTableListChangedNotification)notification).removed, s -> s.descriptor);
else
return Collections.emptyList();
}
private static Map<VersionAndType, Integer> update(Map<VersionAndType, Integer> counts,
VersionAndType toUpdate,
int delta)
{
Map<VersionAndType, Integer> m = counts == null ? new HashMap<>() : counts;
m.merge(toUpdate, delta, (a, b) -> (a + b == 0) ? null : (a + b));
return m;
}
@VisibleForTesting
static VersionAndType version(Descriptor sstable)
{
return new VersionAndType(sstable.version, sstable.formatType);
}
}

View File

@ -0,0 +1,50 @@
/*
* 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.service;
import com.google.common.collect.ImmutableSet;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.apache.cassandra.notifications.INotification;
/**
* Notification triggered by a {@link SSTablesGlobalTracker} when the set of sstables versions in use on this node
* changes.
*
* <p>The notification includes the set of sstable versions in use when the notification is triggered (so the result
* of the change triggering that notification).
*/
public class SSTablesVersionsInUseChangeNotification implements INotification
{
/**
* The set of all sstable versions in use on this node at the time of this notification.
*/
public final ImmutableSet<VersionAndType> versionsInUse;
SSTablesVersionsInUseChangeNotification(ImmutableSet<VersionAndType> versionsInUse)
{
this.versionsInUse = versionsInUse;
}
@Override
public String toString()
{
return String.format("SSTablesInUseChangeNotification(%s)", versionsInUse);
}
}

View File

@ -74,6 +74,8 @@ import org.apache.cassandra.gms.*;
import org.apache.cassandra.hints.HintVerbHandler;
import org.apache.cassandra.hints.HintsService;
import org.apache.cassandra.io.sstable.SSTableLoader;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.locator.*;
import org.apache.cassandra.metrics.StorageMetrics;
@ -250,6 +252,8 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
private final StreamStateStore streamStateStore = new StreamStateStore();
public final SSTablesGlobalTracker sstablesTracker;
public boolean isSurveyMode()
{
return isSurveyMode;
@ -325,6 +329,8 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
MessagingService.instance().registerVerbHandlers(MessagingService.Verb.BATCH_STORE, new BatchStoreVerbHandler());
MessagingService.instance().registerVerbHandlers(MessagingService.Verb.BATCH_REMOVE, new BatchRemoveVerbHandler());
sstablesTracker = new SSTablesGlobalTracker(SSTableFormat.Type.current());
}
public void registerDaemon(CassandraDaemon daemon)
@ -898,6 +904,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
appStates.put(ApplicationState.HOST_ID, valueFactory.hostId(localHostId));
appStates.put(ApplicationState.RPC_ADDRESS, valueFactory.rpcaddress(FBUtilities.getBroadcastRpcAddress()));
appStates.put(ApplicationState.RELEASE_VERSION, valueFactory.releaseVersion());
appStates.put(ApplicationState.SSTABLE_VERSIONS, valueFactory.sstableVersions(sstablesTracker.versionsInUse()));
// load the persisted ring state. This used to be done earlier in the init process,
// but now we always perform a shadow round when preparing to join and we have to
@ -908,6 +915,18 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
Gossiper.instance.register(this);
Gossiper.instance.start(SystemKeyspace.incrementAndGetGeneration(), appStates); // needed for node-ring gathering.
gossipActive = true;
sstablesTracker.register((notification, o) -> {
if (!(notification instanceof SSTablesVersionsInUseChangeNotification))
return;
Set<VersionAndType> versions = ((SSTablesVersionsInUseChangeNotification)notification).versionsInUse;
logger.debug("Updating local sstables version in Gossip to {}", versions);
Gossiper.instance.addLocalApplicationState(ApplicationState.SSTABLE_VERSIONS,
valueFactory.sstableVersions(versions));
});
// gossip snitch infos (local DC and rack)
gossipSnitchInfo();
// gossip Schema.emptyVersion forcing immediate check for schema updates (see MigrationManager#maybeScheduleSchemaPull)

View File

@ -33,6 +33,10 @@ import org.apache.cassandra.distributed.shared.Versions;
import static org.apache.cassandra.distributed.shared.AssertUtils.assertRows;
import static org.apache.cassandra.distributed.shared.AssertUtils.row;
import static org.junit.Assert.assertEquals;
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
{
@ -54,7 +58,7 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
for (int i = 0; i < cluster.size(); i++)
{
int nodeNum = i + 1;
System.out.println(String.format("****** node %s: %s", nodeNum, cluster.get(nodeNum).config()));
System.out.printf("****** node %s: %s%n", nodeNum, cluster.get(nodeNum).config());
}
})
.runAfterNodeUpgrade(((cluster, node) -> {
@ -88,7 +92,7 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
for (int i = 0; i < cluster.size(); i++)
{
int nodeNum = i + 1;
System.out.println(String.format("****** node %s: %s", nodeNum, cluster.get(nodeNum).config()));
System.out.printf("****** node %s: %s%n", nodeNum, cluster.get(nodeNum).config());
}
})
.runAfterNodeUpgrade(((cluster, node) -> {
@ -116,13 +120,13 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
new TestCase()
.nodes(2)
.upgrade(Versions.Major.v22, Versions.Major.v3X)
.withConfig(config -> config.with(GOSSIP, NETWORK, NATIVE_PROTOCOL))
.setup(cluster -> {
cluster.schemaChange(String.format(
"CREATE TABLE %s.%s (key int, c1 int, c2 int, c3 int, PRIMARY KEY (key, c1, c2)) WITH COMPACT STORAGE",
KEYSPACE, table));
ICoordinator coordinator = cluster.coordinator(1);
for (int i = 1; i <= partitions; i++)
{
for (int j = 1; j <= rowsPerPartition; j++)
@ -160,32 +164,33 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
});
}).runAfterClusterUpgrade(cluster ->
{
for (int i = 1; i <= cluster.size(); i++)
{
NodeToolResult result = cluster.get(1).nodetoolResult("upgradesstables");
assertEquals("upgrade sstables failed for node " + i, 0, result.getRc());
}
{
for (int i = 1; i <= cluster.size(); i++)
{
NodeToolResult result = cluster.get(i).nodetoolResult("upgradesstables");
assertEquals("upgrade sstables failed for node " + i, 0, result.getRc());
}
Thread.sleep(1000);
// make sure the results are the same after upgrade and upgrade sstables but before dropping compact storage
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
// make sure the results are the same after upgrade and upgrade sstables but before dropping compact storage
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
// make sure the results are the same after dropping compact storage on only the first node
IMessageFilters.Filter filter = cluster.verbs().allVerbs().to(2).drop();
cluster.schemaChange(String.format("ALTER TABLE %s.%s DROP COMPACT STORAGE", KEYSPACE, table), 1);
// make sure the results are the same after dropping compact storage on only the first node
IMessageFilters.Filter filter = cluster.verbs().allVerbs().to(2).drop();
cluster.schemaChange(String.format("ALTER TABLE %s.%s DROP COMPACT STORAGE", KEYSPACE, table), 1);
recorder.validateResults(cluster, 1, ConsistencyLevel.ONE);
recorder.validateResults(cluster, 1, ConsistencyLevel.ONE);
filter.off();
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
filter.off();
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
// make sure the results continue to be the same after dropping compact storage on the second node
cluster.schemaChange(String.format("ALTER TABLE %s.%s DROP COMPACT STORAGE", KEYSPACE, table), 2);
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
})
// make sure the results continue to be the same after dropping compact storage on the second node
cluster.schemaChange(String.format("ALTER TABLE %s.%s DROP COMPACT STORAGE", KEYSPACE, table), 2);
recorder.validateResults(cluster, 1);
recorder.validateResults(cluster, 2);
})
.run();
}
@ -200,6 +205,7 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
new TestCase()
.nodes(2)
.upgrade(Versions.Major.v22, Versions.Major.v3X)
.withConfig(config -> config.with(GOSSIP, NETWORK, NATIVE_PROTOCOL))
.setup(cluster -> {
cluster.schemaChange(String.format(
"CREATE TABLE %s.%s (key int, c1 int, c2 int, c3 int, PRIMARY KEY (key, c1, c2)) WITH COMPACT STORAGE",
@ -221,11 +227,8 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
})
.runAfterClusterUpgrade(cluster -> {
for (int i = 1; i <= cluster.size(); i++)
{
NodeToolResult result = cluster.get(1).nodetoolResult("upgradesstables");
assertEquals("upgrade sstables failed for node " + i, 0, result.getRc());
}
cluster.forEach(n -> n.nodetoolResult("upgradesstables", KEYSPACE).asserts().success());
Thread.sleep(1000);
// drop compact storage on only one node before performing writes
IMessageFilters.Filter filter = cluster.verbs().allVerbs().to(2).drop();
@ -262,7 +265,7 @@ public class CompactStorage2to3UpgradeTest extends UpgradeTestBase
KEYSPACE, table, 8, 1, 3), ConsistencyLevel.ALL);
coordinator.execute(String.format("DELETE FROM %s.%s WHERE key = %d and c1 = %d and c2 > 1",
KEYSPACE, table, 6, 2, 4), ConsistencyLevel.ALL);
KEYSPACE, table, 6, 2), ConsistencyLevel.ALL);
ResultsRecorder recorder = new ResultsRecorder();
runQueries(coordinator, recorder, new String[] {

View File

@ -0,0 +1,72 @@
/*
* 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.upgrade;
import org.junit.Test;
import org.apache.cassandra.distributed.api.ConsistencyLevel;
import org.apache.cassandra.distributed.api.NodeToolResult;
import org.apache.cassandra.distributed.shared.Versions;
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.catchThrowable;
import static org.junit.Assert.assertEquals;
public class DropCompactStorageTest extends UpgradeTestBase
{
@Test
public void dropCompactStorageBeforeUpgradesstablesTo3X() throws Throwable
{
dropCompactStorageBeforeUpgradeSstables(Versions.Major.v3X);
}
/**
* Upgrades a node from 2.2 to 3.x and DROP COMPACT just after the upgrade but _before_ upgrading the underlying
* sstables.
*
* <p>This test reproduces the issue from CASSANDRA-15897.
*/
public void dropCompactStorageBeforeUpgradeSstables(Versions.Major upgradeTo) throws Throwable
{
new TestCase()
.nodes(1)
.upgrade(Versions.Major.v22, upgradeTo)
.withConfig(config -> config.with(GOSSIP, NETWORK, NATIVE_PROTOCOL))
.setup((cluster) -> {
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (id int, ck int, v int, PRIMARY KEY (id, ck)) WITH COMPACT STORAGE");
for (int i = 0; i < 5; i++)
cluster.coordinator(1).execute("INSERT INTO "+KEYSPACE+".tbl (id, ck, v) values (1, ?, ?)", ConsistencyLevel.ALL, i, i);
cluster.get(1).flush(KEYSPACE);
})
.runAfterNodeUpgrade((cluster, node) -> {
Throwable thrown = catchThrowable(() -> cluster.schemaChange("ALTER TABLE "+KEYSPACE+".tbl DROP COMPACT STORAGE"));
assertThat(thrown).hasMessageContainingAll("Cannot DROP COMPACT STORAGE as some nodes in the cluster",
"has some non-upgraded 2.x sstables");
assertThat(cluster.get(1).nodetool("upgradesstables")).isEqualTo(0);
Thread.sleep(1000);
cluster.schemaChange("ALTER TABLE "+KEYSPACE+".tbl DROP COMPACT STORAGE");
cluster.coordinator(1).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL);
})
.run();
}
}

View File

@ -86,14 +86,14 @@ public class TrackerTest
List<SSTableReader> readers = ImmutableList.of(MockSchema.sstable(0, true, cfs), MockSchema.sstable(1, cfs), MockSchema.sstable(2, cfs));
tracker.addInitialSSTables(copyOf(readers));
Assert.assertNull(tracker.tryModify(ImmutableList.of(MockSchema.sstable(0, cfs)), OperationType.COMPACTION));
try (LifecycleTransaction txn = tracker.tryModify(readers.get(0), OperationType.COMPACTION);)
try (LifecycleTransaction txn = tracker.tryModify(readers.get(0), OperationType.COMPACTION))
{
Assert.assertNotNull(txn);
Assert.assertNull(tracker.tryModify(readers.get(0), OperationType.COMPACTION));
Assert.assertEquals(1, txn.originals().size());
Assert.assertTrue(txn.originals().contains(readers.get(0)));
}
try (LifecycleTransaction txn = tracker.tryModify(Collections.<SSTableReader>emptyList(), OperationType.COMPACTION);)
try (LifecycleTransaction txn = tracker.tryModify(Collections.<SSTableReader>emptyList(), OperationType.COMPACTION))
{
Assert.assertNotNull(txn);
Assert.assertEquals(0, txn.originals().size());
@ -147,12 +147,17 @@ public class TrackerTest
{
ColumnFamilyStore cfs = MockSchema.newCFS();
Tracker tracker = cfs.getTracker();
MockListener listener = new MockListener(false);
tracker.subscribe(listener);
List<SSTableReader> readers = ImmutableList.of(MockSchema.sstable(0, 17, cfs),
MockSchema.sstable(1, 121, cfs),
MockSchema.sstable(2, 9, cfs));
tracker.addInitialSSTables(copyOf(readers));
Assert.assertEquals(3, tracker.view.get().sstables.size());
Assert.assertEquals(1, listener.senders.size());
Assert.assertEquals(1, listener.received.size());
Assert.assertTrue(listener.received.get(0) instanceof InitialSSTableAddedNotification);
for (SSTableReader reader : readers)
Assert.assertTrue(reader.isKeyCacheSetup());
@ -181,6 +186,7 @@ public class TrackerTest
Assert.assertEquals(17 + 121 + 9, cfs.metric.liveDiskSpaceUsed.getCount());
Assert.assertEquals(1, listener.senders.size());
Assert.assertEquals(1, listener.received.size());
Assert.assertEquals(tracker, listener.senders.get(0));
Assert.assertTrue(listener.received.get(0) instanceof SSTableAddedNotification);
DatabaseDescriptor.setIncrementalBackupsEnabled(backups);
@ -236,15 +242,16 @@ public class TrackerTest
Assert.assertNull(tracker.dropSSTables(reader -> reader != readers.get(0), OperationType.UNKNOWN, null));
Assert.assertEquals(1, tracker.getView().sstables.size());
Assert.assertEquals(3, listener.received.size());
Assert.assertEquals(4, listener.received.size());
Assert.assertEquals(tracker, listener.senders.get(0));
Assert.assertTrue(listener.received.get(0) instanceof SSTableDeletingNotification);
Assert.assertTrue(listener.received.get(1) instanceof SSTableDeletingNotification);
Assert.assertTrue(listener.received.get(2) instanceof SSTableListChangedNotification);
Assert.assertEquals(readers.get(1), ((SSTableDeletingNotification) listener.received.get(0)).deleting);
Assert.assertEquals(readers.get(2), ((SSTableDeletingNotification)listener.received.get(1)).deleting);
Assert.assertEquals(2, ((SSTableListChangedNotification) listener.received.get(2)).removed.size());
Assert.assertEquals(0, ((SSTableListChangedNotification) listener.received.get(2)).added.size());
Assert.assertTrue(listener.received.get(0) instanceof InitialSSTableAddedNotification);
Assert.assertTrue(listener.received.get(1) instanceof SSTableDeletingNotification);
Assert.assertTrue(listener.received.get(2) instanceof SSTableDeletingNotification);
Assert.assertTrue(listener.received.get(3) instanceof SSTableListChangedNotification);
Assert.assertEquals(readers.get(1), ((SSTableDeletingNotification) listener.received.get(1)).deleting);
Assert.assertEquals(readers.get(2), ((SSTableDeletingNotification)listener.received.get(2)).deleting);
Assert.assertEquals(2, ((SSTableListChangedNotification) listener.received.get(3)).removed.size());
Assert.assertEquals(0, ((SSTableListChangedNotification) listener.received.get(3)).added.size());
Assert.assertEquals(9, cfs.metric.liveDiskSpaceUsed.getCount());
readers.get(0).selfRef().release();
}
@ -339,7 +346,7 @@ public class TrackerTest
Tracker tracker = new Tracker(null, false);
MockListener listener = new MockListener(false);
tracker.subscribe(listener);
tracker.notifyAdded(singleton(r1));
tracker.notifyAdded(singleton(r1), false);
Assert.assertEquals(singleton(r1), ((SSTableAddedNotification) listener.received.get(0)).added);
listener.received.clear();
tracker.notifyDeleting(r1);
@ -360,7 +367,7 @@ public class TrackerTest
MockListener failListener = new MockListener(true);
tracker.subscribe(failListener);
tracker.subscribe(listener);
Assert.assertNotNull(tracker.notifyAdded(singleton(r1), null));
Assert.assertNotNull(tracker.notifyAdded(singleton(r1), false, null));
Assert.assertEquals(singleton(r1), ((SSTableAddedNotification) listener.received.get(0)).added);
listener.received.clear();
Assert.assertNotNull(tracker.notifySSTablesChanged(singleton(r1), singleton(r2), OperationType.COMPACTION, null));

View File

@ -0,0 +1,51 @@
/*
* 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.io.sstable.format;
import org.junit.Test;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static junit.framework.Assert.assertEquals;
public class VersionAndTypeTest
{
@Test
public void testValidInput()
{
assertEquals(VersionAndType.fromString("big-bc").toString(), "big-bc");
}
@Test
public void testInvalidInputs()
{
assertThatThrownBy(() -> VersionAndType.fromString(" ")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("should be of the form 'big-bc'");
assertThatThrownBy(() -> VersionAndType.fromString("mcc-")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("should be of the form 'big-bc'");
assertThatThrownBy(() -> VersionAndType.fromString("mcc-d")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("No Type constant mcc");
assertThatThrownBy(() -> VersionAndType.fromString("mcc-d-")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("No Type constant mcc");
assertThatThrownBy(() -> VersionAndType.fromString("mcc")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("should be of the form 'big-bc'");
assertThatThrownBy(() -> VersionAndType.fromString("-")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("should be of the form 'big-bc'");
assertThatThrownBy(() -> VersionAndType.fromString("--")).isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("should be of the form 'big-bc'");
}
}

View File

@ -0,0 +1,158 @@
/*
* 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.service;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
import org.junit.Test;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.VersionAndType;
import org.assertj.core.util.Files;
import org.quicktheories.core.Gen;
import org.quicktheories.generators.Generate;
import static org.junit.Assert.assertEquals;
import static org.quicktheories.QuickTheory.qt;
import static org.quicktheories.generators.SourceDSL.integers;
import static org.quicktheories.generators.SourceDSL.lists;
import static org.quicktheories.generators.SourceDSL.strings;
public class SSTablesGlobalTrackerTest
{
private static final int MAX_VERSION_LIST_SIZE = 10;
private static final int MAX_UPDATES_PER_GEN = 100;
/**
* Ensures that the tracker properly maintains the set of versions in use.
*
* <p>Using 'Quick Theories', we generate a number of random sstables notification and validate that after each
* update the tracker computes the proper set of versions in use (using a simplisitic model that keeps the set
* of all sstables in use and maps it to the set of its versions).
*/
@Test
public void testUpdates()
{
qt().forAll(lists().of(updates()).ofSizeBetween(0, MAX_UPDATES_PER_GEN),
sstableFormatTypes())
.checkAssert((updates, formatType) -> {
SSTablesGlobalTracker tracker = new SSTablesGlobalTracker(formatType);
Set<Descriptor> all = new HashSet<>();
Set<VersionAndType> previous = Collections.emptySet();
for (Update update : updates)
{
update.applyTo(all);
boolean triggerUpdate = tracker.handleSSTablesChange(update.removed, update.added);
Set<VersionAndType> expectedInUse = versionAndTypes(all);
assertEquals(expectedInUse, tracker.versionsInUse());
assertEquals(!expectedInUse.equals(previous), triggerUpdate);
previous = expectedInUse;
}
});
}
private Set<VersionAndType> versionAndTypes(Set<Descriptor> descriptors)
{
return descriptors.stream().map(SSTablesGlobalTracker::version).collect(Collectors.toSet());
}
private Gen<String> keyspaces()
{
return Generate.pick(Arrays.asList("k1", "k2"));
}
private Gen<String> tables()
{
return Generate.pick(Arrays.asList("t1", "t2", "t3"));
}
private Gen<Integer> generations()
{
return integers().between(1, 20);
}
private Gen<Descriptor> descriptors()
{
return sstableFormatTypes().zip(keyspaces(),
tables(),
generations(),
sstableVersionString(),
(f, k, t, g, v) -> new Descriptor(v, Files.currentFolder(), k, t, g, f));
}
private Gen<List<Descriptor>> descriptorLists(int minSize)
{
return lists().of(descriptors()).ofSizeBetween(minSize, MAX_VERSION_LIST_SIZE);
}
private Gen<SSTableFormat.Type> sstableFormatTypes()
{
return Generate.enumValues(SSTableFormat.Type.class);
}
private Gen<String> sstableVersionString()
{
// We want to somewhat favor the current version, as that is technically more realistic so we generate it 50%
// of the time, and generate something random 50% of the time.
return Generate.constant(SSTableFormat.Type.current().info.getLatestVersion().getVersion())
.mix(strings().betweenCodePoints('a', 'z').ofLength(2));
}
private Gen<Update> updates()
{
// We want to avoid update that remove and add nothing, as those appear to be generated quite a bit are not
// too useful. Yet, having one of removed/added be empty is actually something we want.
Gen<Update> maybeEmptyRemoved = descriptorLists(0).zip(descriptorLists(1), Update::new);
Gen<Update> maybeEmptyAdded = descriptorLists(1).zip(descriptorLists(0), Update::new);
return maybeEmptyRemoved.mix(maybeEmptyAdded);
}
private static class Update
{
final List<Descriptor> removed;
final List<Descriptor> added;
Update(List<Descriptor> removed, List<Descriptor> added)
{
this.removed = removed;
this.added = added;
}
void applyTo(Set<Descriptor> allSSTables)
{
allSSTables.removeAll(removed);
allSSTables.addAll(added);
}
@Override
public String toString()
{
return "Update{" +
"removed=" + removed +
", added=" + added +
'}';
}
}
}