mirror of https://github.com/apache/cassandra
Fix CQLSSTableWriter serialization of vector of date and time
patch by Lukasz Antoniak; reviewed by Andres de la Pena, Yifan Cai for CASSANDRA-20979
This commit is contained in:
parent
f894b8440d
commit
0136fc9c8f
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -1787,12 +1787,6 @@ public abstract class TypeCodec<T>
|
|||
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<T>
|
|||
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)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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<Integer, ?> valueFactory,
|
||||
BiConsumer<Integer, List<?>> 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<Object> 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);
|
||||
|
|
|
|||
Loading…
Reference in New Issue