Merge branch 'cassandra-4.0' into cassandra-4.1

This commit is contained in:
Stefan Miklosovic 2024-09-15 19:54:41 +02:00
commit 48eab9afd8
No known key found for this signature in database
GPG Key ID: 32F35CB2F546D93E
10 changed files with 432 additions and 139 deletions

View File

@ -6,6 +6,7 @@
* Reduce info logging from automatic paxos repair (CASSANDRA-19445)
* Support legacy plain_text_auth section in credentials file removed unintentionally (CASSANDRA-19498)
Merged from 4.0:
* Fix millisecond and microsecond precision for commit log replay (CASSANDRA-19448)
* Improve accuracy of memtable heap usage tracking (CASSANDRA-17298)
* Fix rendering UNSET collection types in query tracing (CASSANDRA-19880)
* Fix latency reported by ideal consistency level monitoring (CASSANDRA-19651)

View File

@ -37,7 +37,11 @@ restore_command=
restore_directories=
# Restore mutations created up to and including this timestamp in GMT.
# Format: yyyy:MM:dd HH:mm:ss (2012:04:31 20:43:12)
# There are only three different formats to express three time precisions:
# Seconds, Milliseconds, and Microseconds.
# Seconds format: yyyy:MM:dd HH:mm:ss (2012:04:31 20:43:12)
# Milliseconds format: yyyy:MM:dd HH:mm:ss.SSS (2012:04:31 20:43:12.633)
# Microseconds format: yyyy:MM:dd HH:mm:ss.SSSSSS (2012:04:31 20:43:12.633222)
#
# Recovery will continue through the segment when the first client-supplied
# timestamp greater than this time is encountered, but only mutations less than

View File

@ -80,7 +80,7 @@ public class CommitLog implements CommitLogMBean
final public AbstractCommitLogSegmentManager segmentManager;
public final CommitLogArchiver archiver;
public CommitLogArchiver archiver;
public final CommitLogMetrics metrics;
final AbstractCommitLogService executor;
@ -162,7 +162,7 @@ public class CommitLog implements CommitLogMBean
return getUnmanagedFiles().length > 0;
}
private File[] getUnmanagedFiles()
public File[] getUnmanagedFiles()
{
File[] files = new File(segmentManager.storageDirectory).tryList(unmanagedFilesFilter);
if (files == null)
@ -397,7 +397,7 @@ public class CommitLog implements CommitLogMBean
@Override
public long getRestorePointInTime()
{
return archiver.restorePointInTime;
return archiver.restorePointInTimeInMicroseconds;
}
@Override
@ -406,6 +406,12 @@ public class CommitLog implements CommitLogMBean
return archiver.precision.toString();
}
@VisibleForTesting
public void setCommitlogArchiver(CommitLogArchiver archiver)
{
this.archiver = archiver;
}
public List<String> getActiveSegmentNames()
{
Collection<CommitLogSegment> segments = segmentManager.getActiveSegments();

View File

@ -23,14 +23,25 @@ package org.apache.cassandra.db.commitlog;
import java.io.IOException;
import java.io.InputStream;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.time.Instant;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.time.format.DateTimeParseException;
import java.util.Map;
import java.util.Properties;
import java.util.TimeZone;
import java.util.concurrent.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Strings;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.io.util.File;
@ -39,50 +50,47 @@ import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.WrappedRunnable;
import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.google.common.base.Strings;
import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory;
public class CommitLogArchiver
{
private static final Logger logger = LoggerFactory.getLogger(CommitLogArchiver.class);
public static final SimpleDateFormat format = new SimpleDateFormat("yyyy:MM:dd HH:mm:ss");
public static final DateTimeFormatter format = DateTimeFormatter.ofPattern("yyyy:MM:dd HH:mm:ss[.[SSSSSS][SSS]]").withZone(ZoneId.of("GMT"));
private static final String COMMITLOG_ARCHIVNG_PROPERTIES_FILE_NAME = "commitlog_archiving.properties";
private static final String DELIMITER = ",";
private static final Pattern NAME = Pattern.compile("%name");
private static final Pattern PATH = Pattern.compile("%path");
private static final Pattern FROM = Pattern.compile("%from");
private static final Pattern TO = Pattern.compile("%to");
static
{
format.setTimeZone(TimeZone.getTimeZone("GMT"));
}
public final Map<String, Future<?>> archivePending = new ConcurrentHashMap<String, Future<?>>();
public final Map<String, Future<?>> archivePending = new ConcurrentHashMap<>();
private final ExecutorService executor;
final String archiveCommand;
final String restoreCommand;
final String restoreDirectories;
public long restorePointInTime;
public CommitLogPosition snapshotCommitLogPosition;
public final TimeUnit precision;
TimeUnit precision;
long restorePointInTimeInMicroseconds;
final CommitLogPosition snapshotCommitLogPosition;
public CommitLogArchiver(String archiveCommand, String restoreCommand, String restoreDirectories,
long restorePointInTime, CommitLogPosition snapshotCommitLogPosition, TimeUnit precision)
public CommitLogArchiver(String archiveCommand,
String restoreCommand,
String restoreDirectories,
long restorePointInTimeInMicroseconds,
CommitLogPosition snapshotCommitLogPosition,
TimeUnit precision)
{
this.archiveCommand = archiveCommand;
this.restoreCommand = restoreCommand;
this.restoreDirectories = restoreDirectories;
this.restorePointInTime = restorePointInTime;
this.restorePointInTimeInMicroseconds = restorePointInTimeInMicroseconds;
this.snapshotCommitLogPosition = snapshotCommitLogPosition;
this.precision = precision;
executor = !Strings.isNullOrEmpty(archiveCommand)
? executorFactory()
.withJmxInternal()
.sequential("CommitLogArchiver")
: null;
? executorFactory()
.withJmxInternal()
.sequential("CommitLogArchiver")
: null;
}
public static CommitLogArchiver disabled()
@ -92,73 +100,96 @@ public class CommitLogArchiver
public static CommitLogArchiver construct()
{
Properties commitlog_commands = new Properties();
try (InputStream stream = CommitLogArchiver.class.getClassLoader().getResourceAsStream("commitlog_archiving.properties"))
Properties commitlogProperties = new Properties();
try (InputStream stream = CommitLogArchiver.class.getClassLoader().getResourceAsStream(COMMITLOG_ARCHIVNG_PROPERTIES_FILE_NAME))
{
if (stream == null)
{
logger.trace("No commitlog_archiving properties found; archive + pitr will be disabled");
logger.trace("No {} found; archiving and point-in-time-restoration will be disabled", COMMITLOG_ARCHIVNG_PROPERTIES_FILE_NAME);
return disabled();
}
else
{
commitlog_commands.load(stream);
String archiveCommand = commitlog_commands.getProperty("archive_command");
String restoreCommand = commitlog_commands.getProperty("restore_command");
String restoreDirectories = commitlog_commands.getProperty("restore_directories");
if (restoreDirectories != null && !restoreDirectories.isEmpty())
{
for (String dir : restoreDirectories.split(DELIMITER))
{
File directory = new File(dir);
if (!directory.exists())
{
if (!directory.tryCreateDirectory())
{
throw new RuntimeException("Unable to create directory: " + dir);
}
}
}
}
String targetTime = commitlog_commands.getProperty("restore_point_in_time");
TimeUnit precision = TimeUnit.valueOf(commitlog_commands.getProperty("precision", "MICROSECONDS"));
long restorePointInTime;
try
{
restorePointInTime = Strings.isNullOrEmpty(targetTime) ? Long.MAX_VALUE : format.parse(targetTime).getTime();
}
catch (ParseException e)
{
throw new RuntimeException("Unable to parse restore target time", e);
}
String snapshotPosition = commitlog_commands.getProperty("snapshot_commitlog_position");
CommitLogPosition snapshotCommitLogPosition;
try
{
snapshotCommitLogPosition = Strings.isNullOrEmpty(snapshotPosition)
? CommitLogPosition.NONE
: CommitLogPosition.serializer.fromString(snapshotPosition);
}
catch (ParseException | NumberFormatException e)
{
throw new RuntimeException("Unable to parse snapshot commit log position", e);
}
return new CommitLogArchiver(archiveCommand,
restoreCommand,
restoreDirectories,
restorePointInTime,
snapshotCommitLogPosition,
precision);
commitlogProperties.load(stream);
return getArchiverFromProperties(commitlogProperties);
}
}
catch (IOException e)
{
throw new RuntimeException("Unable to load commitlog_archiving.properties", e);
throw new RuntimeException("Unable to load " + COMMITLOG_ARCHIVNG_PROPERTIES_FILE_NAME, e);
}
}
@VisibleForTesting
static CommitLogArchiver getArchiverFromProperties(Properties commitlogCommands)
{
assert !commitlogCommands.isEmpty();
String archiveCommand = commitlogCommands.getProperty("archive_command");
String restoreCommand = commitlogCommands.getProperty("restore_command");
String restoreDirectories = commitlogCommands.getProperty("restore_directories");
if (restoreDirectories != null && !restoreDirectories.isEmpty())
{
for (String dir : restoreDirectories.split(DELIMITER))
{
File directory = new File(dir);
if (!directory.exists())
{
if (!directory.tryCreateDirectory())
{
throw new RuntimeException("Unable to create directory: " + dir);
}
}
}
}
String precisionPropertyValue = commitlogCommands.getProperty("precision", TimeUnit.MICROSECONDS.name());
TimeUnit precision;
try
{
precision = TimeUnit.valueOf(precisionPropertyValue);
}
catch (IllegalArgumentException ex)
{
throw new RuntimeException("Unable to parse precision of value " + precisionPropertyValue, ex);
}
if (precision == TimeUnit.NANOSECONDS)
throw new RuntimeException("NANOSECONDS level precision is not supported.");
String targetTime = commitlogCommands.getProperty("restore_point_in_time");
long restorePointInTime = Long.MAX_VALUE;
try
{
if (!Strings.isNullOrEmpty(targetTime))
{
// get restorePointInTime in microseconds level by default as cassandra use this level's timestamp
restorePointInTime = getRestorationPointInTimeInMicroseconds(targetTime);
}
}
catch (DateTimeParseException e)
{
throw new RuntimeException("Unable to parse restore target time", e);
}
String snapshotPosition = commitlogCommands.getProperty("snapshot_commitlog_position");
CommitLogPosition snapshotCommitLogPosition;
try
{
snapshotCommitLogPosition = Strings.isNullOrEmpty(snapshotPosition)
? CommitLogPosition.NONE
: CommitLogPosition.serializer.fromString(snapshotPosition);
}
catch (ParseException | NumberFormatException e)
{
throw new RuntimeException("Unable to parse snapshot commit log position", e);
}
return new CommitLogArchiver(archiveCommand,
restoreCommand,
restoreDirectories,
restorePointInTime,
snapshotCommitLogPosition,
precision);
}
public void maybeArchive(final CommitLogSegment segment)
@ -182,7 +213,7 @@ public class CommitLogArchiver
* Differs from the above because it can be used on any file, rather than only
* managed commit log segments (and thus cannot call waitForFinalSync), and in
* the treatment of failures.
*
* <p>
* Used to archive files present in the commit log directory at startup (CASSANDRA-6904).
* Since the files being already archived by normal operation could cause subsequent
* hard-linking or other operations to fail, we should not throw errors on failure
@ -192,20 +223,16 @@ public class CommitLogArchiver
if (Strings.isNullOrEmpty(archiveCommand))
return;
archivePending.put(name, executor.submit(new Runnable()
{
public void run()
archivePending.put(name, executor.submit(() -> {
try
{
try
{
String command = NAME.matcher(archiveCommand).replaceAll(Matcher.quoteReplacement(name));
command = PATH.matcher(command).replaceAll(Matcher.quoteReplacement(path));
exec(command);
}
catch (IOException e)
{
logger.warn("Archiving file {} failed, file may have already been archived.", name, e);
}
String command = NAME.matcher(archiveCommand).replaceAll(Matcher.quoteReplacement(name));
command = PATH.matcher(command).replaceAll(Matcher.quoteReplacement(path));
exec(command);
}
catch (IOException e)
{
logger.warn("Archiving file {} failed, file may have already been archived.", name, e);
}
}));
}
@ -230,7 +257,7 @@ public class CommitLogArchiver
{
if (e.getCause().getCause() instanceof IOException)
{
logger.error("Looks like the archiving of file {} failed earlier, cassandra is going to ignore this segment for now.", name, e.getCause().getCause());
logger.error("Looks like the archiving of file {} failed earlier, Cassandra is going to ignore this segment for now.", name, e.getCause().getCause());
return false;
}
}
@ -310,4 +337,36 @@ public class CommitLogArchiver
pb.redirectErrorStream(true);
FBUtilities.exec(pb);
}
}
/**
* We change the restore_point_in_time from configuration file into microseconds level as Cassandra use microseconds
* as the timestamp.
*
* @param restorationPointInTime value of "restore_point_in_time" in properties file.
* @return microseconds value of restore_point_in_time
*/
@VisibleForTesting
public static long getRestorationPointInTimeInMicroseconds(String restorationPointInTime)
{
assert !Strings.isNullOrEmpty(restorationPointInTime) : "restore_point_in_time is null or empty!";
Instant instant = format.parse(restorationPointInTime, Instant::from);
return instant.getEpochSecond() * 1_000_000 + instant.getNano() / 1000;
}
public long getRestorePointInTimeInMicroseconds()
{
return this.restorePointInTimeInMicroseconds;
}
@VisibleForTesting
public void setRestorePointInTimeInMicroseconds(long restorePointInTimeInMicroseconds)
{
this.restorePointInTimeInMicroseconds = restorePointInTimeInMicroseconds;
}
@VisibleForTesting
public void setPrecision(TimeUnit timeUnit)
{
this.precision = timeUnit;
}
}

View File

@ -41,7 +41,11 @@ public interface CommitLogMBean
/**
* Restore mutations created up to and including this timestamp in GMT
* Format: yyyy:MM:dd HH:mm:ss (2012:04:31 20:43:12)
* There are only three different formats to express three time precisions:
* Seconds, Milliseconds, and Microseconds.
* Seconds format: yyyy:MM:dd HH:mm:ss (2012:04:31 20:43:12)
* Milliseconds format: yyyy:MM:dd HH:mm:ss.SSS (2012:04:31 20:43:12.633)
* Microseconds format: yyyy:MM:dd HH:mm:ss.SSSSSS (2012:04:31 20:43:12.633222)
*
* Recovery will continue through the segment when the first client-supplied
* timestamp greater than this time is encountered, but only mutations less than

View File

@ -19,7 +19,16 @@
package org.apache.cassandra.db.commitlog;
import java.io.IOException;
import java.util.*;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Queue;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
import com.google.common.annotations.VisibleForTesting;
@ -28,8 +37,6 @@ import com.google.common.collect.HashMultimap;
import com.google.common.collect.Iterables;
import com.google.common.collect.Multimap;
import com.google.common.collect.Ordering;
import org.apache.cassandra.io.util.File;
import org.apache.commons.lang3.StringUtils;
import org.apache.cassandra.utils.concurrent.Future;
@ -46,6 +53,7 @@ import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.SystemKeyspace;
import org.apache.cassandra.db.partitions.PartitionUpdate;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.schema.Schema;
@ -76,7 +84,7 @@ public class CommitLogReplayer implements CommitLogReadHandler
private long pendingMutationBytes = 0;
private final ReplayFilter replayFilter;
private final CommitLogArchiver archiver;
private CommitLogArchiver archiver;
@VisibleForTesting
protected boolean sawCDCMutation;
@ -115,7 +123,8 @@ public class CommitLogReplayer implements CommitLogReadHandler
// Point in time restore is taken to mean that the tables need to be replayed even if they were
// deleted at a later point in time. Any truncation record after that point must thus be cleared prior
// to replay (CASSANDRA-9195).
long restoreTime = commitLog.archiver.restorePointInTime;
// truncatedTime is millseconds level but restoreTime is microlevel
long restoreTime = commitLog.archiver.restorePointInTimeInMicroseconds == Long.MAX_VALUE ? Long.MAX_VALUE : commitLog.archiver.restorePointInTimeInMicroseconds / 1000;
long truncatedTime = SystemKeyspace.getTruncatedAt(cfs.metadata.id);
if (truncatedTime > restoreTime)
{
@ -141,7 +150,7 @@ public class CommitLogReplayer implements CommitLogReadHandler
}
else
{
if (commitLog.archiver.restorePointInTime == Long.MAX_VALUE)
if (commitLog.archiver.getRestorePointInTimeInMicroseconds() == Long.MAX_VALUE)
{
// Normal restart, everything is persisted and restored by the memtable itself.
filter = new IntervalSet<>(CommitLogPosition.NONE, CommitLog.instance.getCurrentPosition());
@ -455,11 +464,9 @@ public class CommitLogReplayer implements CommitLogReadHandler
protected boolean pointInTimeExceeded(Mutation fm)
{
long restoreTarget = archiver.restorePointInTime;
for (PartitionUpdate upd : fm.getPartitionUpdates())
{
if (archiver.precision.toMillis(upd.maxTimestamp()) > restoreTarget)
if (archiver.precision.toMicros(upd.maxTimestamp()) > archiver.restorePointInTimeInMicroseconds)
return true;
}
return false;

View File

@ -246,7 +246,8 @@ public abstract class CommitLogSegment
/**
* FOR TESTING PURPOSES.
*/
static void resetReplayLimit()
@VisibleForTesting
public static void resetReplayLimit()
{
replayLimitId = getNextId();
}

View File

@ -38,17 +38,17 @@ public class DropRecreateAndRestoreTest extends CQLTester
{
createTable("CREATE TABLE %s (a int, b int, c int, PRIMARY KEY(a, b))");
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 0, 0);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 0, 1, 1);
long timeInMicroSecond1 = System.currentTimeMillis() * 1000;
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 0, 0, 0, timeInMicroSecond1);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ?", 0, 1, 1, timeInMicroSecond1);
long time = System.currentTimeMillis();
TableId id = currentTableMetadata().id;
assertRows(execute("SELECT * FROM %s"), row(0, 0, 0), row(0, 1, 1));
Thread.sleep(5);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 0, 2);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?)", 1, 1, 3);
long timeInMicroSecond2 = System.currentTimeMillis() * 1000;
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 1, 0, 2, timeInMicroSecond2);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 1, 1, 3, timeInMicroSecond2);
assertRows(execute("SELECT * FROM %s"), row(1, 0, 2), row(1, 1, 3), row(0, 0, 0), row(0, 1, 1));
// Drop will flush and clean segments. Hard-link them so that they can be restored later.
@ -68,13 +68,13 @@ public class DropRecreateAndRestoreTest extends CQLTester
FileUtils.renameWithConfirm(new File(logPath, segment + ".save"), new File(logPath, segment));
try
{
// Restore to point in time.
CommitLog.instance.archiver.restorePointInTime = time;
// Restore to point in time (microseconds granularity)
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(timeInMicroSecond1);
CommitLog.instance.resetUnsafe(false);
}
finally
{
CommitLog.instance.archiver.restorePointInTime = Long.MAX_VALUE;
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(Long.MAX_VALUE);
}
assertRows(execute("SELECT * FROM %s"), row(0, 0, 0), row(0, 1, 1));

View File

@ -22,8 +22,10 @@ import java.io.IOException;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.Date;
import java.util.concurrent.*;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicReference;
import org.junit.Assert;
@ -34,9 +36,6 @@ import org.junit.runner.RunWith;
import org.junit.runners.Parameterized;
import org.junit.runners.Parameterized.Parameters;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.SchemaLoader;
import org.apache.cassandra.Util;
import org.apache.cassandra.concurrent.NamedThreadFactory;
@ -46,7 +45,8 @@ import org.apache.cassandra.db.commitlog.CommitLog;
import org.apache.cassandra.db.commitlog.CommitLogArchiver;
import org.apache.cassandra.db.commitlog.CommitLogReplayer;
import org.apache.cassandra.db.context.CounterContext;
import org.apache.cassandra.db.rows.*;
import org.apache.cassandra.db.rows.Row;
import org.apache.cassandra.db.rows.UnfilteredRowIterator;
import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.io.compress.DeflateCompressor;
import org.apache.cassandra.io.compress.LZ4Compressor;
@ -64,8 +64,6 @@ import static org.junit.Assert.assertEquals;
@RunWith(Parameterized.class)
public class RecoveryManagerTest
{
private static Logger logger = LoggerFactory.getLogger(RecoveryManagerTest.class);
private static final String KEYSPACE1 = "RecoveryManagerTest1";
private static final String CF_STANDARD1 = "Standard1";
private static final String CF_STATIC1 = "Static1";
@ -249,14 +247,17 @@ public class RecoveryManagerTest
public void testRecoverPIT() throws Exception
{
CommitLog.instance.resetUnsafe(true);
long originalPIT = CommitLog.instance.archiver.getRestorePointInTimeInMicroseconds();
ColumnFamilyStore cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1);
Date date = CommitLogArchiver.format.parse("2112:12:12 12:12:12");
long timeMS = date.getTime() - 5000;
// seconds level
// the archiver's restorePointInTime use the commitlog_archiving_properties file's
long rpiTs = CommitLogArchiver.getRestorationPointInTimeInMicroseconds("2112:12:12 12:12:12");
long timeInMicroLevel = rpiTs - 5000;
Keyspace keyspace1 = Keyspace.open(KEYSPACE1);
for (int i = 0; i < 10; ++i)
{
long ts = TimeUnit.MILLISECONDS.toMicros(timeMS + (i * 1000));
long ts = timeInMicroLevel + (i * 1000);
new RowUpdateBuilder(cfs.metadata(), ts, "name-" + i)
.clustering("cc")
.add("val", Integer.toString(i))
@ -266,11 +267,66 @@ public class RecoveryManagerTest
// Sanity check row count prior to clear and replay
assertEquals(10, Util.getAll(Util.cmd(cfs).build()).size());
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
CommitLog.instance.resetUnsafe(false);
assertEquals(6, Util.getAll(Util.cmd(cfs).build()).size());
//reset the rpi
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(originalPIT);
keyspace1.getColumnFamilyStore("Standard1").clearUnsafe();
CommitLog.instance.resetUnsafe(false);
CommitLog.instance.resetUnsafe(true);
cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1);
keyspace1 = Keyspace.open(KEYSPACE1);
// milseconds level
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
rpiTs = CommitLogArchiver.getRestorationPointInTimeInMicroseconds("2112:12:12 12:12:12.063");
timeInMicroLevel = rpiTs - 5000;
keyspace1 = Keyspace.open(KEYSPACE1);
for (int i = 0; i < 10; ++i)
{
long ts = timeInMicroLevel + (i * 1000);
new RowUpdateBuilder(cfs.metadata(), ts, "name-" + i)
.clustering("cc")
.add("val", Integer.toString(i))
.build()
.apply();
}
// Sanity check row count prior to clear and replay
assertEquals(10, Util.getAll(Util.cmd(cfs).build()).size());
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(rpiTs);
CommitLog.instance.resetUnsafe(false);
assertEquals(6, Util.getAll(Util.cmd(cfs).build()).size());
//reset the rpi
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(originalPIT);
CommitLog.instance.resetUnsafe(true);
cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1);
keyspace1 = Keyspace.open(KEYSPACE1);
// milseconds level
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
rpiTs = CommitLogArchiver.getRestorationPointInTimeInMicroseconds("2112:12:12 12:12:12.063222");
timeInMicroLevel = rpiTs - 5000;
keyspace1 = Keyspace.open(KEYSPACE1);
for (int i = 0; i < 10; ++i)
{
long ts = timeInMicroLevel + (i * 1000);
new RowUpdateBuilder(cfs.metadata(), ts, "name-" + i)
.clustering("cc")
.add("val", Integer.toString(i))
.build()
.apply();
}
// Sanity check row count prior to clear and replay
assertEquals(10, Util.getAll(Util.cmd(cfs).build()).size());
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(rpiTs);
CommitLog.instance.resetUnsafe(false);
assertEquals(6, Util.getAll(Util.cmd(cfs).build()).size());
//reset the rpi
CommitLog.instance.archiver.setRestorePointInTimeInMicroseconds(originalPIT);
}
@Test
@ -279,13 +335,11 @@ public class RecoveryManagerTest
CommitLog.instance.resetUnsafe(true);
Keyspace keyspace1 = Keyspace.open(KEYSPACE1);
ColumnFamilyStore cfs = keyspace1.getColumnFamilyStore(CF_STATIC1);
Date date = CommitLogArchiver.format.parse("2112:12:12 12:12:12");
long timeMS = date.getTime() - 5000;
long timeInMicroLevel = CommitLogArchiver.getRestorationPointInTimeInMicroseconds("2112:12:12 12:12:12") - 5000;
for (int i = 0; i < 10; ++i)
{
long ts = TimeUnit.MILLISECONDS.toMicros(timeMS + (i * 1000));
long ts = timeInMicroLevel + (i * 1000);
new RowUpdateBuilder(cfs.metadata(), ts, "name-" + i)
.add("val", Integer.toString(i))
.build()
@ -306,8 +360,8 @@ public class RecoveryManagerTest
{
CommitLog.instance.resetUnsafe(true);
ColumnFamilyStore cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1);
Date date = CommitLogArchiver.format.parse("2112:12:12 12:12:12");
long timeMS = date.getTime();
// 2112:12:12 12:12:12 is from the commitlog_archiving.properties file for testing
long timeInMicroLevel = CommitLogArchiver.getRestorationPointInTimeInMicroseconds("2112:12:12 12:12:12");
Keyspace keyspace1 = Keyspace.open(KEYSPACE1);
@ -316,9 +370,9 @@ public class RecoveryManagerTest
{
long ts;
if (i == 9)
ts = TimeUnit.MILLISECONDS.toMicros(timeMS - 1000);
ts = timeInMicroLevel - 1000;
else
ts = TimeUnit.MILLISECONDS.toMicros(timeMS + (i * 1000));
ts = timeInMicroLevel + (i * 1000);
new RowUpdateBuilder(cfs.metadata(), ts, "name-" + i)
.clustering("cc")
@ -330,7 +384,7 @@ public class RecoveryManagerTest
// Sanity check row count prior to clear and replay
assertEquals(10, Util.getAll(Util.cmd(cfs).build()).size());
keyspace1.getColumnFamilyStore("Standard1").clearUnsafe();
keyspace1.getColumnFamilyStore(CF_STANDARD1).clearUnsafe();
CommitLog.instance.resetUnsafe(false);
assertEquals(2, Util.getAll(Util.cmd(cfs).build()).size());

View File

@ -0,0 +1,157 @@
/*
* 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.db.commitlog;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.file.Path;
import java.util.HashMap;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.ClassRule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;
import org.apache.cassandra.cql3.CQLTester;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.RowUpdateBuilder;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.PathUtils;
import static org.apache.cassandra.io.util.PathUtils.forEach;
import static org.junit.Assert.assertTrue;
public class CommitLogArchiverTest extends CQLTester
{
private static Path dirName;
private static final String rpiTime = "2024:03:22 20:43:12.633222";
private File dir;
@ClassRule
public static TemporaryFolder temporaryFolder = new TemporaryFolder();
@BeforeClass
public static void beforeClass() throws IOException
{
dirName = temporaryFolder.newFolder().toPath();
CommitLog commitLog = CommitLog.instance;
Properties properties = new Properties();
properties.putAll(new HashMap<String, String>()
{{
put("archive_command", "/bin/cp %path " + dirName);
put("restore_command", "/bin/cp -f %from %to");
put("restore_directories", dirName.toString());
put("restore_point_in_time", rpiTime);
}});
CommitLogArchiver commitLogArchiver = CommitLogArchiver.getArchiverFromProperties(properties);
commitLog.setCommitlogArchiver(commitLogArchiver);
}
@Before
public void before()
{
dir = new File(dirName);
// to prevent other test cases' archive files from affecting us
if (dir.isDirectory() && dir.tryList().length > 0)
forEach(dirName, PathUtils::deleteRecursive);
}
@Test
public void testArchiver()
{
String table = createTable(KEYSPACE, "CREATE TABLE %s (a TEXT PRIMARY KEY, b blob);");
ColumnFamilyStore cfs = Keyspace.open(KEYSPACE).getColumnFamilyStore(table);
long ts = CommitLogArchiver.getRestorationPointInTimeInMicroseconds(rpiTime);
ByteBuffer value = ByteBuffer.allocate(1024);
// Make sure that new CommitLogSegment will be allocated as the CommitLogSegment size is 5M
// and if new CommitLogSegment is allocated then the old CommitLogSegment will be archived.
for (int i = 1; i <= 10; ++i)
{
new RowUpdateBuilder(cfs.metadata(), ts - i, "name-" + i)
.add("b", value)
.build()
.apply();
}
CommitLog.instance.forceRecycleAllSegments();
CommitLog.instance.segmentManager.awaitManagementTasksCompletion();
// If the number of files that under backup dir is bigger than 1, that means the
// arhiver for commitlog is effective.
assertTrue(dir.isDirectory() && dir.tryList().length > 0);
}
@Test
public void testRestoreInDifferentPrecision() throws Throwable
{
createTable(KEYSPACE, "CREATE TABLE %s (a INT , b INT, c INT, PRIMARY KEY(a, b));");
// default level is microsecond
long timeInMicroSecond1 = CommitLogArchiver.getRestorationPointInTimeInMicroseconds(rpiTime);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 3, 0, 0, timeInMicroSecond1);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 3, 1, 1, timeInMicroSecond1);
long timeInMicroSecond2 = timeInMicroSecond1 + 1;
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ?", 4, 0, 0, timeInMicroSecond2);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ?", 4, 1, 1, timeInMicroSecond2);
assertRows(execute("SELECT * FROM %s"), row(4, 0, 0), row(4, 1, 1), row(3, 0, 0), row(3, 1, 1));
CommitLog.instance.forceRecycleAllSegments();
CommitLog.instance.segmentManager.awaitManagementTasksCompletion();
execute("TRUNCATE TABLE %s");
assertRowCount(execute("SELECT * FROM %s"), 0);
// replay log
CommitLog.instance.archiver.maybeRestoreArchive();
CommitLogSegment.resetReplayLimit();
// restore archived files
CommitLog.instance.recoverFiles(CommitLog.instance.getUnmanagedFiles());
// restore poin time is rpiTime in microseconds , so row(4, 0, 0) and row(4, 1, 1) is skipped
assertRows(execute("SELECT * FROM %s"), row(3, 0, 0), row(3, 1, 1));
// set to millisecond level
CommitLog.instance.archiver.setPrecision(TimeUnit.MILLISECONDS);
long timeInMilliSecond3 = timeInMicroSecond1 / 1000 + 1;
execute("TRUNCATE TABLE %s");
assertRowCount(execute("SELECT * FROM %s"), 0);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 2, 0, 0, timeInMilliSecond3);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 2, 1, 1, timeInMilliSecond3);
long timeInMilliSecond4 = timeInMilliSecond3 - 1;
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 1, 0, 0, timeInMilliSecond4);
execute("INSERT INTO %s (a, b, c) VALUES (?, ?, ?) USING TIMESTAMP ? ", 1, 1, 1, timeInMilliSecond4);
assertRows(execute("SELECT * FROM %s"), row(1, 0, 0), row(1, 1, 1), row(2, 0, 0), row(2, 1, 1));
CommitLog.instance.forceRecycleAllSegments();
CommitLog.instance.segmentManager.awaitManagementTasksCompletion();
execute("TRUNCATE TABLE %s");
assertRowCount(execute("SELECT * FROM %s"), 0);
// replay log
CommitLog.instance.archiver.maybeRestoreArchive();
CommitLogSegment.resetReplayLimit();
CommitLog.instance.recoverFiles(CommitLog.instance.getUnmanagedFiles());
// restore poin time is rpiTime in millseconds, so row(2, 0, 0) and row(2, 1, 1) is skipped
assertRows(execute("SELECT * FROM %s"), row(1, 0, 0), row(1, 1, 1));
}
}