From 2aeb6a02905a39fb5b26c93dccd53da67d0ba24d Mon Sep 17 00:00:00 2001 From: Doug Rohrer Date: Tue, 29 Apr 2025 17:47:56 -0400 Subject: [PATCH] CASSANDRA-20609 - CQLSSTableWriter should support setting the format (5.0) --- CHANGES.txt | 1 + .../io/sstable/CQLSSTableWriter.java | 12 ++++++ .../io/sstable/CQLSSTableWriterTest.java | 39 +++++++++++++++++-- 3 files changed, 48 insertions(+), 4 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 1791454d17..08e0ecfa16 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.0.5 + * CQLSSTableWriter supports setting the format (BTI or Big) (CASSANDRA-20609) * Don't allocate in ThreadLocalReadAheadBuffer#close() (CASSANDRA-20551) * Ensure RowFilter#isMutableIntersection() properly evaluates numeric ranges on a single column (CASSANDRA-20566) * Switch memtable-related off-heap objects to Native Endian and Memory to Little Endian (CASSANDRA-20190) diff --git a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java index 05b23c49f2..64834792fc 100644 --- a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java @@ -636,6 +636,18 @@ public class CQLSSTableWriter implements Closeable return this; } + /** + * Specifies the SSTable format this CQLSSTableWriter instance should use for writing. + * + * @param format The format to use + * @return this builder + */ + public Builder withFormat(SSTableFormat format) + { + this.format = format; + return this; + } + /** * Whether the produced sstable should be open or not. * By default, the writer does not open the produced sstables diff --git a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java index d358185732..cc41b0ea2a 100644 --- a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java @@ -20,6 +20,8 @@ package org.apache.cassandra.io.sstable; import java.io.IOException; import java.nio.ByteBuffer; +import java.nio.file.Files; +import java.nio.file.Path; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -34,6 +36,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiPredicate; import java.util.stream.Collectors; +import java.util.stream.Stream; import java.util.stream.StreamSupport; import com.google.common.collect.ImmutableList; @@ -61,8 +64,10 @@ import org.apache.cassandra.dht.Murmur3Partitioner; import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.index.sai.disk.format.IndexDescriptor; import org.apache.cassandra.index.sai.utils.IndexIdentifier; +import org.apache.cassandra.io.sstable.format.SSTableFormat; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.sstable.format.big.BigFormat; +import org.apache.cassandra.io.sstable.format.bti.BtiFormat; import org.apache.cassandra.io.util.File; import org.apache.cassandra.io.util.PathUtils; import org.apache.cassandra.schema.Schema; @@ -102,9 +107,22 @@ public abstract class CQLSSTableWriterTest } @Test - public void testUnsortedWriter() throws Exception + public void testUnsortedWriterBig() throws Exception { - try (AutoCloseable switcher = Util.switchPartitioner(ByteOrderedPartitioner.instance)) + BigFormat format = BigFormat.getInstance(); + testWritingSstableWithFormat(format); + } + + @Test + public void testUnsortedWriterBti() throws Exception + { + SSTableFormat btiFormat = new BtiFormat.BtiFormatFactory().getInstance(Collections.emptyMap()); + testWritingSstableWithFormat(btiFormat); + } + + private void testWritingSstableWithFormat(SSTableFormat format) throws Exception + { + try (AutoCloseable ignored = Util.switchPartitioner(ByteOrderedPartitioner.instance)) { String schema = "CREATE TABLE " + qualifiedTable + " (" + " k int PRIMARY KEY," @@ -115,6 +133,7 @@ public abstract class CQLSSTableWriterTest CQLSSTableWriter writer = CQLSSTableWriter.builder() .inDirectory(dataDir) .forTable(schema) + .withFormat(format) .using(insert).build(); writer.addRow(0, "test1", 24); @@ -124,6 +143,7 @@ public abstract class CQLSSTableWriterTest writer.close(); + validateFilesAreInFormat(format); loadSSTables(dataDir, keyspace); UntypedResultSet rs = QueryProcessor.executeInternal("SELECT * FROM " + qualifiedTable); @@ -140,7 +160,6 @@ public abstract class CQLSSTableWriterTest row = iter.next(); assertEquals(1, row.getInt("k")); assertEquals("test2", row.getString("v1")); - //assertFalse(row.has("v2")); assertEquals(44, row.getInt("v2")); row = iter.next(); @@ -150,11 +169,23 @@ public abstract class CQLSSTableWriterTest row = iter.next(); assertEquals(3, row.getInt("k")); - assertEquals(null, row.getBytes("v1")); // Using getBytes because we know it won't NPE + assertFalse(row.has("v1")); assertEquals(12, row.getInt("v2")); } } + private void validateFilesAreInFormat(SSTableFormat format) throws IOException + { + try (Stream dataFilePaths = Files.list(dataDir.toPath()).filter(p -> p.toString().endsWith("Data.db"))) + { + dataFilePaths.forEach(dataFilePath -> { + File dataFile = new File(dataFilePath.toFile()); + Descriptor descriptor = Descriptor.fromFile(dataFile); + assertEquals(format, descriptor.version.format); + }); + } + } + @Test public void testForbidCounterUpdates() throws Exception {