From 0136fc9c8fb50418cb5458cab35df341062c52b1 Mon Sep 17 00:00:00 2001 From: Lukasz Antoniak Date: Wed, 17 Dec 2025 15:25:33 +0100 Subject: [PATCH] Fix CQLSSTableWriter serialization of vector of date and time patch by Lukasz Antoniak; reviewed by Andres de la Pena, Yifan Cai for CASSANDRA-20979 --- CHANGES.txt | 1 + .../cql3/functions/types/TypeCodec.java | 13 ++-- .../io/sstable/CQLSSTableWriterTest.java | 77 +++++++++++++++++++ 3 files changed, 85 insertions(+), 6 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 56dc843ab8..2569b03124 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.0.7 + * Fix CQLSSTableWriter serialization of vector of date and time (CASSANDRA-20979) * Correctly calculate default for FailureDetector max interval (CASSANDRA-21025) * Adding missing configs in system_views.settings to be backward compatible (CASSANDRA-20863) * Heap dump should not be generated on handled exceptions (CASSANDRA-20974) diff --git a/src/java/org/apache/cassandra/cql3/functions/types/TypeCodec.java b/src/java/org/apache/cassandra/cql3/functions/types/TypeCodec.java index 1212e8e2ab..9b2df91fbe 100644 --- a/src/java/org/apache/cassandra/cql3/functions/types/TypeCodec.java +++ b/src/java/org/apache/cassandra/cql3/functions/types/TypeCodec.java @@ -1787,12 +1787,6 @@ public abstract class TypeCodec super(DataType.date(), LocalDate.class); } - @Override - public int serializedSize() - { - return 8; - } - @Override public LocalDate parse(String value) { @@ -1876,6 +1870,13 @@ public abstract class TypeCodec super(DataType.time()); } + @Override + public int serializedSize() + { + // matching behavior of TimeType, which is not declared as fixed length + return VARIABLE_LENGTH; + } + @Override public Long parse(String value) { diff --git a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java index cc41b0ea2a..9c88b12ebd 100644 --- a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java @@ -34,7 +34,9 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BiConsumer; import java.util.function.BiPredicate; +import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.Stream; import java.util.stream.StreamSupport; @@ -49,6 +51,7 @@ import org.junit.rules.TemporaryFolder; import com.datastax.driver.core.utils.UUIDs; import org.apache.cassandra.Util; +import org.apache.cassandra.cql3.CQL3Type; import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.cql3.UntypedResultSet; import org.apache.cassandra.cql3.functions.types.DataType; @@ -58,6 +61,10 @@ import org.apache.cassandra.cql3.functions.types.UDTValue; import org.apache.cassandra.cql3.functions.types.UserType; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Keyspace; +import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.FloatType; +import org.apache.cassandra.db.marshal.SimpleDateType; +import org.apache.cassandra.db.marshal.TimeType; import org.apache.cassandra.db.marshal.UTF8Type; import org.apache.cassandra.dht.ByteOrderedPartitioner; import org.apache.cassandra.dht.Murmur3Partitioner; @@ -76,6 +83,7 @@ import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.JavaDriverUtils; import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis; +import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@ -1520,6 +1528,75 @@ public abstract class CQLSSTableWriterTest assertFalse(indexDescriptor.isPerColumnIndexBuildComplete(new IndexIdentifier(keyspace, table, "idx2"))); } + @Test + public void testWritingVectorData() throws Exception + { + testWritingVectorData(CQL3Type.Native.FLOAT, FloatType.instance, (i) -> (float) i, (i, vector) -> { + assertThat(vector).allMatch(val -> val instanceof Float); + assertThat(vector).allMatch(val -> (float) val == (float) i); + }); + + perTestSetup(); + + testWritingVectorData(CQL3Type.Native.DATE, SimpleDateType.instance, LocalDate::fromDaysSinceEpoch, (i, vector) -> { + assertThat(vector).allMatch(val -> val instanceof Integer); + assertThat(vector).allMatch(val -> { + int days = (int) val - Integer.MIN_VALUE; // signed to unsigned conversion + return days == i; + }); + }); + + perTestSetup(); + + testWritingVectorData(CQL3Type.Native.TIME, TimeType.instance, (i) -> (long) i, (i, vector) -> { + assertThat(vector).allMatch(val -> val instanceof Long); + assertThat(vector).allMatch(val -> (long) val == (long) i); + }); + } + + private void testWritingVectorData(CQL3Type.Native cqlType, AbstractType subType, Function valueFactory, + BiConsumer> checkFunction) throws Exception + { + final int dimensions = 5; + final String schema = "CREATE TABLE " + qualifiedTable + " (" + + " k int," + + " v1 VECTOR<" + cqlType.name() + ", " + dimensions + ">," + + " PRIMARY KEY (k)" + + ")"; + + CQLSSTableWriter writer = CQLSSTableWriter.builder() + .inDirectory(dataDir) + .forTable(schema) + .using("INSERT INTO " + keyspace + "." + table + " (k, v1) " + + "VALUES (?, ?)").build(); + + for (int i = 0; i < 100; i++) + { + List vector = new ArrayList<>(dimensions); + for (int j = 0; j < dimensions; j++) + { + vector.add(valueFactory.apply(i)); + } + writer.addRow(i, vector); + } + + writer.close(); + loadSSTables(dataDir, keyspace); + + UntypedResultSet resultSet = QueryProcessor.executeInternal("SELECT * FROM " + keyspace + "." + table); + + assertEquals(resultSet.size(), 100); + int cnt = 0; + for (UntypedResultSet.Row row : resultSet) + { + assertEquals(cnt, row.getInt("k")); + List vector = row.getVector("v1", subType, dimensions); + assertThat(vector).hasSize(dimensions); + checkFunction.accept(cnt, vector); + cnt++; + } + } + protected void loadSSTables(File dataDir, String ksName) { ColumnFamilyStore cfs = Keyspace.openWithoutSSTables(ksName).getColumnFamilyStore(table);