Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -54,18 +54,20 @@ public abstract class AbstractSSTableSimpleWriter implements Closeable
protected final RegularAndStaticColumns columns;
protected SSTableFormat<?, ?> format = DatabaseDescriptor.getSelectedSSTableFormat();
protected static final AtomicReference<SSTableId> id = new AtomicReference<>(SSTableIdFactory.instance.defaultBuilder().generator(Stream.empty()).get());
protected final SSTableId.Builder<? extends SSTableId> idBuilder;
protected boolean makeRangeAware = false;
protected final Collection<Index.Group> indexGroups;
protected Consumer<Collection<SSTableReader>> sstableProducedListener;
protected boolean openSSTableOnProduced = false;
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<? extends SSTableId> idBuilder)
{
this.metadata = metadata;
this.directory = directory;
this.columns = columns;
this.idBuilder = idBuilder;
indexGroups = new ArrayList<>();
}

Expand Down Expand Up @@ -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)
{
Expand All @@ -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;
}
Expand Down
12 changes: 10 additions & 2 deletions src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java
Original file line number Diff line number Diff line change
Expand Up @@ -419,6 +419,7 @@ public static class Builder
private Consumer<Collection<SSTableReader>> sstableProducedListener;
private boolean openSSTableOnProduced = false;
private CompressionDictionary compressionDictionary = null;
private SSTableId.Builder<? extends SSTableId> idBuilder = null;

protected Builder()
{
Expand Down Expand Up @@ -671,6 +672,12 @@ public Builder withFormat(SSTableFormat<?, ?> format)
return this;
}

public Builder withSSTableIdBuilder(SSTableId.Builder<? extends SSTableId> idBuilder)
{
this.idBuilder = idBuilder;
return this;
}

/**
* Use specific compression dictionary upon writing the data.
*
Expand Down Expand Up @@ -799,9 +806,10 @@ public CQLSSTableWriter build()
ModificationStatement preparedModificationStatement = prepareModificationStatement();

TableMetadataRef ref = tableMetadata.ref;
SSTableId.Builder<? extends SSTableId> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<? extends SSTableId> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
}

/**
Expand All @@ -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<? extends SSTableId> idBuilder)
{
super(directory, metadata, columns, idBuilder);
this.maxSSTableSizeInBytes = maxSSTableSizeInMiB * 1024L * 1024L;
this.owner = owner;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down