This is an automated email from the ASF dual-hosted git repository.
smiklosovic pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra.git
The following commit(s) were added to refs/heads/trunk by this push:
new e86ae7f make creation timestamp consistent for all tables of a
snapshot
e86ae7f is described below
commit e86ae7fbe52c95747d69f7d78b6da6fbd34bd48d
Author: Stefan Miklosovic <[email protected]>
AuthorDate: Wed Sep 8 15:40:55 2021 +0200
make creation timestamp consistent for all tables of a snapshot
patch by Stefan Miklosovic; reviewed by Paulo Motta and Aleksei Zotov for
CASSANDRA-16920
---
.../org/apache/cassandra/db/ColumnFamilyStore.java | 38 +++++++++++++++-------
src/java/org/apache/cassandra/db/Keyspace.java | 7 ++--
.../org/apache/cassandra/db/SystemKeyspace.java | 5 ++-
.../cassandra/db/compaction/CompactionTask.java | 7 +++-
.../apache/cassandra/service/StorageService.java | 8 +++--
.../service/snapshot/SnapshotManifest.java | 4 +--
.../{SnapshotsTTLTest.java => SnapshotsTest.java} | 31 ++++++++++++++----
.../org/apache/cassandra/db/DirectoriesTest.java | 4 +--
.../service/snapshot/SnapshotManifestTest.java | 2 +-
9 files changed, 74 insertions(+), 32 deletions(-)
diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
index 05161d7..2cfb3b8 100644
--- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
+++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
@@ -24,6 +24,7 @@ import java.lang.reflect.Constructor;
import java.lang.reflect.InvocationTargetException;
import java.nio.ByteBuffer;
import java.nio.file.Files;
+import java.time.Instant;
import java.util.*;
import java.util.Objects;
import java.util.concurrent.*;
@@ -1524,7 +1525,11 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
{
// skip snapshot creation during scrub, SEE JIRA 5891
if(!disableSnapshot)
- snapshotWithoutFlush("pre-scrub-" + System.currentTimeMillis());
+ {
+ Instant creationTime = Instant.now();
+ String snapshotName = "pre-scrub-" + creationTime.toEpochMilli();
+ snapshotWithoutFlush(snapshotName, creationTime);
+ }
try
{
@@ -1839,13 +1844,18 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
public TableSnapshot snapshotWithoutFlush(String snapshotName)
{
- return snapshotWithoutFlush(snapshotName, null, false, null, null);
+ return snapshotWithoutFlush(snapshotName, Instant.now());
+ }
+
+ public TableSnapshot snapshotWithoutFlush(String snapshotName, Instant
creationTime)
+ {
+ return snapshotWithoutFlush(snapshotName, null, false, null, null,
creationTime);
}
/**
* @param ephemeral If this flag is set to true, the snapshot will be
cleaned during next startup
*/
- public TableSnapshot snapshotWithoutFlush(String snapshotName,
Predicate<SSTableReader> predicate, boolean ephemeral, Duration ttl,
RateLimiter rateLimiter)
+ public TableSnapshot snapshotWithoutFlush(String snapshotName,
Predicate<SSTableReader> predicate, boolean ephemeral, Duration ttl,
RateLimiter rateLimiter, Instant creationTime)
{
if (ephemeral && ttl != null)
{
@@ -1873,17 +1883,17 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
}
}
- return createSnapshot(snapshotName, ephemeral, ttl,
snapshottedSSTables);
+ return createSnapshot(snapshotName, ephemeral, ttl,
snapshottedSSTables, creationTime);
}
- protected TableSnapshot createSnapshot(String tag, boolean ephemeral,
Duration ttl, Set<SSTableReader> sstables) {
+ protected TableSnapshot createSnapshot(String tag, boolean ephemeral,
Duration ttl, Set<SSTableReader> sstables, Instant creationTime) {
Set<File> snapshotDirs = sstables.stream()
.map(s ->
Directories.getSnapshotDirectory(s.descriptor, tag).getAbsoluteFile())
.filter(dir ->
!Directories.isSecondaryIndexFolder(dir)) // Remove secondary index subdirectory
.collect(Collectors.toCollection(HashSet::new));
// Create and write snapshot manifest
- SnapshotManifest manifest = new
SnapshotManifest(mapToDataFilenames(sstables), ttl);
+ SnapshotManifest manifest = new
SnapshotManifest(mapToDataFilenames(sstables), ttl, creationTime);
File manifestFile = getDirectories().getSnapshotManifestFile(tag);
writeSnapshotManifest(manifest, manifestFile);
snapshotDirs.add(manifestFile.getParentFile().getAbsoluteFile()); //
manifest may create empty snapshot dir
@@ -2030,7 +2040,7 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
*/
public TableSnapshot snapshot(String snapshotName)
{
- return snapshot(snapshotName, false, null, null);
+ return snapshot(snapshotName, false, null, null, Instant.now());
}
/**
@@ -2038,11 +2048,13 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
*
* @param snapshotName the name of the associated with the snapshot
* @param skipFlush Skip blocking flush of memtable
+ * @param ttl duration after which the taken snapshot is removed
automatically, if supplied with null, it will never be automatically removed
* @param rateLimiter Rate limiter for hardlinks-per-second
+ * @param creationTime time when this snapshot was taken
*/
- public TableSnapshot snapshot(String snapshotName, boolean skipFlush,
Duration ttl, RateLimiter rateLimiter)
+ public TableSnapshot snapshot(String snapshotName, boolean skipFlush,
Duration ttl, RateLimiter rateLimiter, Instant creationTime)
{
- return snapshot(snapshotName, null, false, skipFlush, ttl,
rateLimiter);
+ return snapshot(snapshotName, null, false, skipFlush, ttl,
rateLimiter, creationTime);
}
@@ -2052,21 +2064,23 @@ public class ColumnFamilyStore implements
ColumnFamilyStoreMBean
*/
public TableSnapshot snapshot(String snapshotName,
Predicate<SSTableReader> predicate, boolean ephemeral, boolean skipFlush)
{
- return snapshot(snapshotName, predicate, ephemeral, skipFlush, null,
null);
+ return snapshot(snapshotName, predicate, ephemeral, skipFlush, null,
null, Instant.now());
}
/**
* @param ephemeral If this flag is set to true, the snapshot will be
cleaned up during next startup
* @param skipFlush Skip blocking flush of memtable
+ * @param ttl duration after which the taken snapshot is removed
automatically, if supplied with null, it will never be automatically removed
* @param rateLimiter Rate limiter for hardlinks-per-second
+ * @param creationTime time when this snapshot was taken
*/
- public TableSnapshot snapshot(String snapshotName,
Predicate<SSTableReader> predicate, boolean ephemeral, boolean skipFlush,
Duration ttl, RateLimiter rateLimiter)
+ public TableSnapshot snapshot(String snapshotName,
Predicate<SSTableReader> predicate, boolean ephemeral, boolean skipFlush,
Duration ttl, RateLimiter rateLimiter, Instant creationTime)
{
if (!skipFlush)
{
forceBlockingFlush();
}
- return snapshotWithoutFlush(snapshotName, predicate, ephemeral, ttl,
rateLimiter);
+ return snapshotWithoutFlush(snapshotName, predicate, ephemeral, ttl,
rateLimiter, creationTime);
}
public boolean snapshotExists(String snapshotName)
diff --git a/src/java/org/apache/cassandra/db/Keyspace.java
b/src/java/org/apache/cassandra/db/Keyspace.java
index 66eb1e0..e5cf741 100644
--- a/src/java/org/apache/cassandra/db/Keyspace.java
+++ b/src/java/org/apache/cassandra/db/Keyspace.java
@@ -19,6 +19,7 @@ package org.apache.cassandra.db;
import java.io.File;
import java.io.IOException;
+import java.time.Instant;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
@@ -248,7 +249,7 @@ public class Keyspace
* @param rateLimiter Rate limiter for hardlinks-per-second
* @throws IOException if the column family doesn't exist
*/
- public void snapshot(String snapshotName, String columnFamilyName, boolean
skipFlush, Duration ttl, RateLimiter rateLimiter) throws IOException
+ public void snapshot(String snapshotName, String columnFamilyName, boolean
skipFlush, Duration ttl, RateLimiter rateLimiter, Instant creationTime) throws
IOException
{
assert snapshotName != null;
boolean tookSnapShot = false;
@@ -257,7 +258,7 @@ public class Keyspace
if (columnFamilyName == null ||
cfStore.name.equals(columnFamilyName))
{
tookSnapShot = true;
- cfStore.snapshot(snapshotName, skipFlush, ttl, rateLimiter);
+ cfStore.snapshot(snapshotName, skipFlush, ttl, rateLimiter,
creationTime);
}
}
@@ -275,7 +276,7 @@ public class Keyspace
*/
public void snapshot(String snapshotName, String columnFamilyName) throws
IOException
{
- snapshot(snapshotName, columnFamilyName, false, null, null);
+ snapshot(snapshotName, columnFamilyName, false, null, null,
Instant.now());
}
/**
diff --git a/src/java/org/apache/cassandra/db/SystemKeyspace.java
b/src/java/org/apache/cassandra/db/SystemKeyspace.java
index 3067142..e07101c 100644
--- a/src/java/org/apache/cassandra/db/SystemKeyspace.java
+++ b/src/java/org/apache/cassandra/db/SystemKeyspace.java
@@ -22,6 +22,7 @@ import java.io.IOError;
import java.io.IOException;
import java.net.InetAddress;
import java.nio.ByteBuffer;
+import java.time.Instant;
import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@@ -1473,8 +1474,10 @@ public final class SystemKeyspace
String snapshotName =
Keyspace.getTimestampedSnapshotName(format("upgrade-%s-%s",
previous,
next));
+
+ Instant creationTime = Instant.now();
for (String keyspace : SchemaConstants.LOCAL_SYSTEM_KEYSPACE_NAMES)
- Keyspace.open(keyspace).snapshot(snapshotName, null, false,
null, null);
+ Keyspace.open(keyspace).snapshot(snapshotName, null, false,
null, null, creationTime);
}
}
diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java
b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java
index 13c9725..19a0698 100644
--- a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java
+++ b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java
@@ -17,6 +17,7 @@
*/
package org.apache.cassandra.db.compaction;
+import java.time.Instant;
import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
@@ -115,7 +116,11 @@ public class CompactionTask extends AbstractCompactionTask
CompactionStrategyManager strategy =
cfs.getCompactionStrategyManager();
if (DatabaseDescriptor.isSnapshotBeforeCompaction())
- cfs.snapshotWithoutFlush(System.currentTimeMillis() + "-compact-"
+ cfs.name);
+ {
+ Instant creationTime = Instant.now();
+ cfs.snapshotWithoutFlush(creationTime.toEpochMilli() + "-compact-"
+ cfs.name, creationTime);
+ }
+
try (CompactionController controller =
getCompactionController(transaction.originals()))
{
diff --git a/src/java/org/apache/cassandra/service/StorageService.java
b/src/java/org/apache/cassandra/service/StorageService.java
index a976a28..dcd2598 100644
--- a/src/java/org/apache/cassandra/service/StorageService.java
+++ b/src/java/org/apache/cassandra/service/StorageService.java
@@ -22,6 +22,7 @@ import java.net.InetAddress;
import java.net.UnknownHostException;
import java.nio.ByteBuffer;
import java.nio.file.Paths;
+import java.time.Instant;
import java.util.*;
import java.util.Map.Entry;
import java.util.concurrent.*;
@@ -3873,10 +3874,11 @@ public class StorageService extends
NotificationBroadcasterSupport implements IE
RateLimiter snapshotRateLimiter =
DatabaseDescriptor.getSnapshotRateLimiter();
+ Instant creationTime = Instant.now();
for (Keyspace keyspace : keyspaces)
{
- keyspace.snapshot(tag, null, skipFlush, ttl, snapshotRateLimiter);
+ keyspace.snapshot(tag, null, skipFlush, ttl, snapshotRateLimiter,
creationTime);
}
}
@@ -3938,13 +3940,13 @@ public class StorageService extends
NotificationBroadcasterSupport implements IE
}
RateLimiter snapshotRateLimiter =
DatabaseDescriptor.getSnapshotRateLimiter();
+ Instant creationTime = Instant.now();
for (Entry<Keyspace, List<String>> entry :
keyspaceColumnfamily.entrySet())
{
for (String table : entry.getValue())
- entry.getKey().snapshot(tag, table, skipFlush, ttl,
snapshotRateLimiter);
+ entry.getKey().snapshot(tag, table, skipFlush, ttl,
snapshotRateLimiter, creationTime);
}
-
}
private void verifyKeyspaceIsValid(String keyspaceName)
diff --git
a/src/java/org/apache/cassandra/service/snapshot/SnapshotManifest.java
b/src/java/org/apache/cassandra/service/snapshot/SnapshotManifest.java
index 8fbc619..5d44acb 100644
--- a/src/java/org/apache/cassandra/service/snapshot/SnapshotManifest.java
+++ b/src/java/org/apache/cassandra/service/snapshot/SnapshotManifest.java
@@ -60,10 +60,10 @@ public class SnapshotManifest
this.expiresAt = null;
}
- public SnapshotManifest(List<String> files, Duration ttl)
+ public SnapshotManifest(List<String> files, Duration ttl, Instant
creationTime)
{
this.files = files;
- this.createdAt = Instant.now();
+ this.createdAt = creationTime;
this.expiresAt = ttl == null ? null :
createdAt.plusMillis(ttl.toMilliseconds());
}
diff --git
a/test/distributed/org/apache/cassandra/distributed/test/SnapshotsTTLTest.java
b/test/distributed/org/apache/cassandra/distributed/test/SnapshotsTest.java
similarity index 88%
rename from
test/distributed/org/apache/cassandra/distributed/test/SnapshotsTTLTest.java
rename to
test/distributed/org/apache/cassandra/distributed/test/SnapshotsTest.java
index ace6c11..c56adfe 100644
---
a/test/distributed/org/apache/cassandra/distributed/test/SnapshotsTTLTest.java
+++ b/test/distributed/org/apache/cassandra/distributed/test/SnapshotsTest.java
@@ -19,8 +19,10 @@
package org.apache.cassandra.distributed.test;
import java.io.IOException;
+import java.util.Arrays;
import org.junit.AfterClass;
+import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@@ -32,10 +34,9 @@ import
org.apache.cassandra.distributed.api.IInvokableInstance;
import org.apache.cassandra.distributed.api.NodeToolResult;
import org.apache.cassandra.distributed.shared.WithProperties;
-import static java.lang.String.format;
import static
org.apache.cassandra.distributed.shared.ClusterUtils.stopUnchecked;
-public class SnapshotsTTLTest extends TestBaseImpl
+public class SnapshotsTest extends TestBaseImpl
{
public static final Integer SNAPSHOT_CLEANUP_PERIOD_SECONDS = 1;
public static final Integer FIVE_SECONDS = 5;
@@ -157,7 +158,8 @@ public class SnapshotsTTLTest extends TestBaseImpl
}
@Test
- public void testSecondaryIndexCleanup() throws Exception {
+ public void testSecondaryIndexCleanup() throws Exception
+ {
cluster.schemaChange("CREATE KEYSPACE IF NOT EXISTS default WITH
replication = {'class': 'SimpleStrategy', 'replication_factor': 2};");
cluster.schemaChange("CREATE TABLE default.tbl (key int, value text,
PRIMARY KEY (key))");
cluster.schemaChange("CREATE INDEX value_idx ON default.tbl (value)");
@@ -180,10 +182,25 @@ public class SnapshotsTTLTest extends TestBaseImpl
listSnapshotsResult.stdoutNotContains("first");
}
- private void populate(Cluster cluster) {
- for (int i = 0; i < 100; i++) {
- cluster.coordinator(1).execute("INSERT INTO default.tbl (key,
value) VALUES (?, 'txt')", ConsistencyLevel.ONE, i);
- }
+ @Test
+ public void testSameTimestampOnEachTableOfSnaphot()
+ {
+ cluster.get(1).nodetoolResult("snapshot", "-t",
"sametimestamp").asserts().success();
+ NodeToolResult result = cluster.get(1).nodetoolResult("listsnapshots");
+
+ long distinctTimestamps = Arrays.stream(result.getStdout().split("\n"))
+ .filter(line ->
line.startsWith("sametimestamp"))
+ .map(line -> line.replaceAll(" +", "
").split(" ")[7])
+ .distinct()
+ .count();
+ // assert all dates are same so there is just one value accross all
individual tables
+ Assert.assertEquals(1, distinctTimestamps);
+ }
+
+ private void populate(Cluster cluster)
+ {
+ for (int i = 0; i < 100; i++)
+ cluster.coordinator(1).execute("INSERT INTO default.tbl (key,
value) VALUES (?, 'txt')", ConsistencyLevel.ONE, i);
}
}
diff --git a/test/unit/org/apache/cassandra/db/DirectoriesTest.java
b/test/unit/org/apache/cassandra/db/DirectoriesTest.java
index fbecc69..4c282a6 100644
--- a/test/unit/org/apache/cassandra/db/DirectoriesTest.java
+++ b/test/unit/org/apache/cassandra/db/DirectoriesTest.java
@@ -185,7 +185,7 @@ public class DirectoriesTest
if (createManifest)
{
File manifestFile =
Directories.getSnapshotManifestFile(snapshotDir);
- manifest = new
SnapshotManifest(Collections.singletonList(sstableDesc.filenameFor(Component.DATA)),
new Duration("1m"));
+ manifest = new
SnapshotManifest(Collections.singletonList(sstableDesc.filenameFor(Component.DATA)),
new Duration("1m"), Instant.now());
manifest.serializeToJsonFile(manifestFile);
}
@@ -310,7 +310,7 @@ public class DirectoriesTest
File manifestFile = directories.getSnapshotManifestFile(tag);
- SnapshotManifest manifest = new SnapshotManifest(files, new
Duration("1m"));
+ SnapshotManifest manifest = new SnapshotManifest(files, new
Duration("1m"), Instant.now());
manifest.serializeToJsonFile(manifestFile);
Set<File> dirs = new HashSet<>();
diff --git
a/test/unit/org/apache/cassandra/service/snapshot/SnapshotManifestTest.java
b/test/unit/org/apache/cassandra/service/snapshot/SnapshotManifestTest.java
index 8760270..2dc73bf 100644
--- a/test/unit/org/apache/cassandra/service/snapshot/SnapshotManifestTest.java
+++ b/test/unit/org/apache/cassandra/service/snapshot/SnapshotManifestTest.java
@@ -106,7 +106,7 @@ public class SnapshotManifestTest
@Test
public void testSerializeAndDeserialize() throws Exception {
- SnapshotManifest manifest = new SnapshotManifest(Arrays.asList("db1",
"db2", "db3"), new Duration("2m"));
+ SnapshotManifest manifest = new SnapshotManifest(Arrays.asList("db1",
"db2", "db3"), new Duration("2m"), Instant.now());
File manifestFile = tempFolder.newFile("manifest.json");
manifest.serializeToJsonFile(manifestFile);
manifest = SnapshotManifest.deserializeFromJsonFile(manifestFile);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]