diff --git a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableSimpleWriter.java b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableSimpleWriter.java index 2487e66fd280..db843618d8d2 100644 --- a/src/java/org/apache/cassandra/io/sstable/AbstractSSTableSimpleWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/AbstractSSTableSimpleWriter.java @@ -54,6 +54,7 @@ public abstract class AbstractSSTableSimpleWriter implements Closeable protected final RegularAndStaticColumns columns; protected SSTableFormat format = DatabaseDescriptor.getSelectedSSTableFormat(); protected static final AtomicReference id = new AtomicReference<>(SSTableIdFactory.instance.defaultBuilder().generator(Stream.empty()).get()); + protected final SSTableId.Builder idBuilder; protected boolean makeRangeAware = false; protected final Collection indexGroups; protected Consumer> sstableProducedListener; @@ -61,11 +62,12 @@ public abstract class AbstractSSTableSimpleWriter implements Closeable protected CompressionDictionary compressionDictionary; protected SSTable.Owner owner; - protected AbstractSSTableSimpleWriter(File directory, TableMetadataRef metadata, RegularAndStaticColumns columns) + protected AbstractSSTableSimpleWriter(File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, SSTableId.Builder idBuilder) { this.metadata = metadata; this.directory = directory; this.columns = columns; + this.idBuilder = idBuilder; indexGroups = new ArrayList<>(); } @@ -147,13 +149,13 @@ protected SSTableTxnWriter createWriter(SSTable.Owner owner) throws IOException effectiveOwner); } - private static Descriptor createDescriptor(File directory, final String keyspace, final String columnFamily, final SSTableFormat fmt) throws IOException + private Descriptor createDescriptor(File directory, final String keyspace, final String columnFamily, final SSTableFormat fmt) throws IOException { SSTableId nextGen = getNextId(directory, columnFamily); return new Descriptor(directory, keyspace, columnFamily, nextGen, fmt); } - private static SSTableId getNextId(File directory, final String columnFamily) throws IOException + private SSTableId getNextId(File directory, final String columnFamily) throws IOException { while (true) { @@ -165,7 +167,7 @@ private static SSTableId getNextId(File directory, final String columnFamily) th .map(d -> d.id); SSTableId lastId = id.get(); - SSTableId newId = SSTableIdFactory.instance.defaultBuilder().generator(Stream.concat(existingIds, Stream.of(lastId))).get(); + SSTableId newId = idBuilder.generator(Stream.concat(existingIds, Stream.of(lastId))).get(); if (id.compareAndSet(lastId, newId)) return newId; } diff --git a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java index 71f6efb161cc..d3efcca24487 100644 --- a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java @@ -419,6 +419,7 @@ public static class Builder private Consumer> sstableProducedListener; private boolean openSSTableOnProduced = false; private CompressionDictionary compressionDictionary = null; + private SSTableId.Builder idBuilder = null; protected Builder() { @@ -671,6 +672,12 @@ public Builder withFormat(SSTableFormat format) return this; } + public Builder withSSTableIdBuilder(SSTableId.Builder idBuilder) + { + this.idBuilder = idBuilder; + return this; + } + /** * Use specific compression dictionary upon writing the data. * @@ -799,9 +806,10 @@ public CQLSSTableWriter build() ModificationStatement preparedModificationStatement = prepareModificationStatement(); TableMetadataRef ref = tableMetadata.ref; + SSTableId.Builder effectiveIdBuilder = idBuilder != null ? idBuilder : SSTableIdFactory.instance.defaultBuilder(); AbstractSSTableSimpleWriter writer = sorted - ? new SSTableSimpleWriter(cfs, directory, ref, preparedModificationStatement.updatedColumns(), maxSSTableSizeInMiB) - : new SSTableSimpleUnsortedWriter(cfs, directory, ref, preparedModificationStatement.updatedColumns(), maxSSTableSizeInMiB); + ? new SSTableSimpleWriter(cfs, directory, ref, preparedModificationStatement.updatedColumns(), maxSSTableSizeInMiB, effectiveIdBuilder) + : new SSTableSimpleUnsortedWriter(cfs, directory, ref, preparedModificationStatement.updatedColumns(), maxSSTableSizeInMiB, effectiveIdBuilder); if (format != null) writer.setSSTableFormatType(format); diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java index d1c6a5fbc0d9..938c3b8452ee 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java @@ -72,12 +72,17 @@ class SSTableSimpleUnsortedWriter extends AbstractSSTableSimpleWriter public SSTableSimpleUnsortedWriter(File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB) { - this(null, directory, metadata, columns, maxSSTableSizeInMiB); + this(null, directory, metadata, columns, maxSSTableSizeInMiB, SSTableIdFactory.instance.defaultBuilder()); } SSTableSimpleUnsortedWriter(SSTable.Owner owner, File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB) { - super(directory, metadata, columns); + this(owner, directory, metadata, columns, maxSSTableSizeInMiB, SSTableIdFactory.instance.defaultBuilder()); + } + + SSTableSimpleUnsortedWriter(SSTable.Owner owner, File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB, SSTableId.Builder idBuilder) + { + super(directory, metadata, columns, idBuilder); this.maxSStableSizeInBytes = maxSSTableSizeInMiB * 1024L * 1024L; this.header = new SerializationHeader(true, metadata.get(), columns, EncodingStats.NO_STATS); this.helper = new SerializationHelper(this.header); diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleWriter.java index 7b5fb8a89e45..dc2ed1a58ddd 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleWriter.java @@ -66,7 +66,7 @@ class SSTableSimpleWriter extends AbstractSSTableSimpleWriter */ public SSTableSimpleWriter(File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB) { - this(null, directory, metadata, columns, maxSSTableSizeInMiB); + this(null, directory, metadata, columns, maxSSTableSizeInMiB, SSTableIdFactory.instance.defaultBuilder()); } /** @@ -83,7 +83,12 @@ public SSTableSimpleWriter(File directory, TableMetadataRef metadata, RegularAnd */ protected SSTableSimpleWriter(SSTable.Owner owner, File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB) { - super(directory, metadata, columns); + this(owner, directory, metadata, columns, maxSSTableSizeInMiB, SSTableIdFactory.instance.defaultBuilder()); + } + + protected SSTableSimpleWriter(SSTable.Owner owner, File directory, TableMetadataRef metadata, RegularAndStaticColumns columns, long maxSSTableSizeInMiB, SSTableId.Builder idBuilder) + { + super(directory, metadata, columns, idBuilder); this.maxSSTableSizeInBytes = maxSSTableSizeInMiB * 1024L * 1024L; this.owner = owner; } diff --git a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java index a6b4cc1b8777..2cdecea5268b 100644 --- a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java @@ -146,6 +146,53 @@ public void testUnsortedWriterBti() throws Exception testWritingSstableWithFormat(btiFormat); } + @Test + public void testUUIDSSTableIdentifiers() throws Exception + { + String schema = "CREATE TABLE " + qualifiedTable + " (" + + " k int PRIMARY KEY," + + " v int" + + ")"; + String insert = "INSERT INTO " + qualifiedTable + " (k, v) VALUES (?, ?)"; + CQLSSTableWriter writer = CQLSSTableWriter.builder() + .inDirectory(dataDir) + .forTable(schema) + .withSSTableIdBuilder(UUIDBasedSSTableId.Builder.instance) + .using(insert).build(); + + writer.addRow(0, 1); + writer.close(); + + File[] dataFiles = dataDir.tryList((dir, name) -> name.contains("Data.db")); + assertEquals(1, dataFiles.length); + + Descriptor descriptor = Descriptor.fromFile(dataFiles[0]); + assertTrue(UUIDBasedSSTableId.Builder.instance.isUniqueIdentifier(descriptor.id.toString())); + } + + @Test + public void testDefaultSSTableIdentifiersAreLegacy() throws Exception + { + String schema = "CREATE TABLE " + qualifiedTable + " (" + + " k int PRIMARY KEY," + + " v int" + + ")"; + String insert = "INSERT INTO " + qualifiedTable + " (k, v) VALUES (?, ?)"; + CQLSSTableWriter writer = CQLSSTableWriter.builder() + .inDirectory(dataDir) + .forTable(schema) + .using(insert).build(); + + writer.addRow(0, 1); + writer.close(); + + File[] dataFiles = dataDir.tryList((dir, name) -> name.contains("Data.db")); + assertEquals(1, dataFiles.length); + + Descriptor descriptor = Descriptor.fromFile(dataFiles[0]); + assertTrue(SequenceBasedSSTableId.Builder.instance.isUniqueIdentifier(descriptor.id.toString())); + } + private void testWritingSstableWithFormat(SSTableFormat format) throws Exception { try (AutoCloseable ignored = Util.switchPartitioner(ByteOrderedPartitioner.instance))