This is an automated email from the ASF dual-hosted git repository.
lukasz-antoniak pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-analytics.git
The following commit(s) were added to refs/heads/trunk by this push:
new 39ce0567 CASSANALYTICS-26: Support vector data type Patch by Lukasz
Antoniak; reviewed by Shailaja Koppu, Yifan Cai for CASSANALYTICS-26
39ce0567 is described below
commit 39ce05672cb371bccc479d8e8c85d468fe5944dc
Author: Lukasz Antoniak <[email protected]>
AuthorDate: Tue Oct 14 14:09:46 2025 +0200
CASSANALYTICS-26: Support vector data type
Patch by Lukasz Antoniak; reviewed by Shailaja Koppu, Yifan Cai for
CASSANALYTICS-26
---
.circleci/config.yml | 4 +-
.github/workflows/test.yaml | 10 +-
CHANGES.txt | 1 +
build.gradle | 2 +-
cassandra-analytics-cdc/build.gradle | 2 +-
.../java/org/apache/cassandra/cdc/CdcTests.java | 48 ++
.../cassandra/cdc/test/TestVersionSupplier.java | 2 +-
.../apache/cassandra/bridge/CassandraVersion.java | 6 +-
.../cassandra/spark/data/CassandraTypes.java | 11 +
.../org/apache/cassandra/spark/data/CqlField.java | 8 +-
.../org/apache/cassandra/spark/utils/CqlUtils.java | 2 +-
.../spark/bulkwriter/SqlToCqlTypeConverter.java | 2 +
.../cassandra/spark/KryoSerializationTests.java | 27 ++
.../spark/bulkwriter/MockBulkWriterContext.java | 2 +-
.../spark/bulkwriter/RecordWriterTest.java | 2 +-
.../bulkwriter/StreamSessionConsistencyTest.java | 2 +-
.../cassandra/spark/endtoend/DataTypeTests.java | 101 ++++
.../apache/cassandra/spark/endtoend/MiscTests.java | 16 +-
.../build.gradle | 4 +-
.../distributed/impl/CassandraCluster.java | 5 +
.../testing/SharedClusterIntegrationTestBase.java | 41 +-
.../cassandra/analytics/BulkReaderVectorTest.java | 115 +++++
.../cassandra/analytics/BulkWriteVectorTest.java | 115 +++++
.../SparkSqlTypeConverterImplementation.java | 4 +
.../apache/cassandra/bridge/CassandraBridge.java | 5 +
.../data/converter/types/VectorTypeTests.java | 59 +++
.../bridge/CassandraTypesImplementation.java | 19 +
.../cassandra/spark/data/complex/CqlVector.java | 166 +++++++
.../cassandra/spark/reader/SchemaBuilder.java | 114 +++++
.../spark/reader/AbstractStreamScanner.java | 2 +-
.../cassandra/bridge/AbstractCassandraTypes.java | 6 +
.../spark/data/complex/CqlCollection.java | 17 +-
...hemaBuilder.java => AbstractSchemaBuilder.java} | 73 ++-
.../cassandra/spark/reader/SchemaBuilder.java | 516 +--------------------
gradle.properties | 2 +-
gradlew | 6 +
scripts/build-dtest-jars.sh | 4 +-
scripts/relocate-dtest-dependencies.pom | 1 +
38 files changed, 932 insertions(+), 590 deletions(-)
diff --git a/.circleci/config.yml b/.circleci/config.yml
index b3238f77..68e5c11d 100644
--- a/.circleci/config.yml
+++ b/.circleci/config.yml
@@ -374,7 +374,7 @@ workflows:
spark: ["3"]
scala: ["2.13"]
jdk: ["11"]
- cassandra: ["5.0.5"]
+ cassandra: ["5.0.7"]
# Cassandra 5.0 on Spark 4 / Scala 2.13 / JDK 17
- int-test:
@@ -386,4 +386,4 @@ workflows:
spark: ["4"]
scala: ["2.13"]
jdk: ["17"]
- cassandra: ["5.0.5"]
+ cassandra: ["5.0.7"]
diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml
index db7bf223..2bc40c45 100644
--- a/.github/workflows/test.yaml
+++ b/.github/workflows/test.yaml
@@ -223,13 +223,13 @@ jobs:
# into each match. To add a new version: add one entry to 'config' and
one to
# 'include'.
matrix:
- config: ['s2.13-c5.0.5', 's2.12-c4.1.4', 's2.12-c4.0.17',
's2.13-c5.0.5-spark4']
+ config: ['s2.13-c5.0.7', 's2.12-c4.1.4', 's2.12-c4.0.17',
's2.13-c5.0.7-spark4']
job_index: [0, 1, 2, 3, 4]
job_total: [5]
include:
- - config: 's2.13-c5.0.5'
+ - config: 's2.13-c5.0.7'
scala: '2.13'
- cassandra: '5.0.5'
+ cassandra: '5.0.7'
jdk: '11'
spark: '3'
- config: 's2.12-c4.1.4'
@@ -242,9 +242,9 @@ jobs:
cassandra: '4.0.17'
jdk: '11'
spark: '3'
- - config: 's2.13-c5.0.5-spark4'
+ - config: 's2.13-c5.0.7-spark4'
scala: '2.13'
- cassandra: '5.0.5'
+ cassandra: '5.0.7'
jdk: '17'
spark: '4'
fail-fast: false
diff --git a/CHANGES.txt b/CHANGES.txt
index afd381b8..7834ba38 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
0.5.0
-----
+ * Support vector data type (CASSANALYTICS-26)
* CDC batch-write mixing a CDC-enabled and CDC-disabled table drops the CDC
table's mutation (CASSANALYTICS-182)
* CdcState.ReplicaCountSerializer map-size overflow corrupts persisted CDC
state (CASSANALYTICS-184)
* SSTable-version-based bridge determination (CASSANALYTICS-24)
diff --git a/build.gradle b/build.gradle
index 8d9f39ed..7551f06e 100644
--- a/build.gradle
+++ b/build.gradle
@@ -67,7 +67,7 @@ ext.dependencyLocation = (System.getenv("CASSANDRA_DEP_DIR")
?: "${rootDir}/depe
// - cassandraFullVersionMap values must match the supported_versions default
// NOTE: Both maps must ALSO stay in sync with the values in
build-dtest-jars.sh
ext.cassandraVersionEnumMap = ["4.0": "FOURZERO", "4.1": "FOURONE", "5.0":
"FIVEZERO"]
-ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0": "5.0.5"]
+ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0": "5.0.7"]
// Shared helper: sets implemented_versions and supported_versions system
properties on a Test task.
// When majorMinor is provided (e.g. "4.0"), uses that version directly.
diff --git a/cassandra-analytics-cdc/build.gradle
b/cassandra-analytics-cdc/build.gradle
index fb39be34..5a3918cc 100644
--- a/cassandra-analytics-cdc/build.gradle
+++ b/cassandra-analytics-cdc/build.gradle
@@ -159,7 +159,7 @@ def configureCdcTestTask = { Test task, String majorMinor =
null ->
// Full version format to match CDC's TestVersionSupplier; tests both
versions for backward compat.
// 4.1 intentionally excluded from gradlew defaults to keep local
iteration fast;
// use testCassandra41 for targeted 4.1 runs. CI covers 4.1 via
CASSANDRA_VERSION env var.
- task.systemProperty "cassandra.sidecar.versions_to_test",
"4.0.17,5.0.5"
+ task.systemProperty "cassandra.sidecar.versions_to_test",
"4.0.17,5.0.7"
}
task.minHeapSize = '1024m'
diff --git
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
index b4490cb0..d24365a5 100644
---
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
+++
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
@@ -49,6 +49,7 @@ import java.util.stream.Stream;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSet;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
+import org.apache.commons.lang3.StringUtils;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
import org.slf4j.Logger;
@@ -84,6 +85,8 @@ import org.apache.cassandra.spark.data.CqlTable;
import org.apache.cassandra.spark.data.ReplicationFactor;
import org.apache.cassandra.spark.data.partitioner.CassandraInstance;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
+import org.apache.cassandra.spark.data.types.Duration;
+import org.apache.cassandra.spark.data.types.TimeUUID;
import org.apache.cassandra.spark.utils.AsyncExecutor;
import org.apache.cassandra.spark.utils.ByteBufferUtils;
import org.apache.cassandra.spark.utils.IOUtils;
@@ -104,6 +107,7 @@ import static
org.apache.cassandra.cdc.test.CdcTester.testWith;
import static org.apache.cassandra.spark.CommonTestUtils.cql3Type;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.fail;
+import static org.assertj.core.api.Assumptions.assumeThat;
import static org.quicktheories.QuickTheory.qt;
import static org.quicktheories.generators.SourceDSL.arbitrary;
@@ -624,6 +628,50 @@ public class CdcTests extends CdcTestBase
.run());
}
+ @ParameterizedTest
+
@MethodSource("org.apache.cassandra.cdc.test.TestVersionSupplier#testVersions")
+ public void testVector(CassandraVersion version)
+ {
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().forAll(cql3Type(bridge))
+ // Cassandra VectorType does not support swapping custom subtype
serializer,
+ // so we cannot use AnalyticsTimeUUIDSerializer or
AnalyticsDurationSerializer.
+ .assuming(t -> !t.cqlName().equals(Duration.INSTANCE.name()) &&
!t.cqlName().equals(TimeUUID.INSTANCE.name()))
+ .checkAssert(
+ t ->
+ testWith(bridge, cdcBridge, commitLogDir,
TestSchema.builder(bridge)
+
.withPartitionKey("pk", bridge.uuid())
+
.withColumn("c1", bridge.bigint())
+
.withColumn("c2", bridge.vector(t, 5)))
+ .withCdcEventChecker((testRows, events) -> {
+ for (CdcEvent event : events)
+ {
+ assertThat(event.getPartitionKeys().size()).isEqualTo(1);
+
assertThat(event.getPartitionKeys().get(0).columnName).isEqualTo("pk");
+ assertThat(event.getClusteringKeys()).isNull();
+ assertThat(event.getStaticColumns()).isNull();
+ assertThat(event.getValueColumns().stream()
+ .map(v -> v.columnName)
+
.collect(Collectors.toList())).isEqualTo(Arrays.asList("c1", "c2"));
+ Value vectorValue = event.getValueColumns().get(1);
+ String vectorType = vectorValue.columnType;
+ assertThat(vectorType.startsWith("vector<")).isTrue();
+ assertThat(vectorType.endsWith(">")).isTrue();
+ assertCqlTypeEquals(t.cqlName(),
+
vectorType.substring(vectorType.indexOf("<") + 1, vectorType.indexOf(","))); //
extract the type in vector<?, ?>
+ String dimensions = StringUtils.substringAfter(vectorType,
",");
+ dimensions = dimensions.substring(0, dimensions.length() -
1).trim();
+ assertThat(dimensions).isEqualTo("5");
+ Object v =
bridge.parseType(vectorType).deserializeToJavaType(vectorValue.getValue());
+ assertThat(v).isInstanceOf(List.class);
+ List list = (List) v;
+ assertThat(list.size()).isGreaterThan(0);
+ assertThat(event.getTtl()).isNull();
+ }
+ })
+ .run());
+ }
+
@ParameterizedTest
@MethodSource("org.apache.cassandra.cdc.test.TestVersionSupplier#testVersions")
public void testMap(CassandraVersion version)
diff --git
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/test/TestVersionSupplier.java
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/test/TestVersionSupplier.java
index 5443e5bf..e6034ffd 100644
---
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/test/TestVersionSupplier.java
+++
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/test/TestVersionSupplier.java
@@ -32,7 +32,7 @@ public final class TestVersionSupplier
public static Stream<CassandraVersion> testVersions()
{
- String versions =
System.getProperty("cassandra.sidecar.versions_to_test", "4.0.17,5.0.5");
+ String versions =
System.getProperty("cassandra.sidecar.versions_to_test", "4.0.17,5.0.7");
return Arrays.stream(versions.split(","))
.map(String::trim)
.map(v -> CassandraVersion.fromVersion(v).orElseThrow(()
-> new IllegalArgumentException("Unsupported version: " + v)));
diff --git
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/bridge/CassandraVersion.java
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/bridge/CassandraVersion.java
index 6ac4b38c..c37116e4 100644
---
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/bridge/CassandraVersion.java
+++
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/bridge/CassandraVersion.java
@@ -38,12 +38,12 @@ import com.google.common.base.Preconditions;
* NOTE: The following values need to stay in sync with:
* - build.gradle:
* - ext.cassandraVersionEnumMap = ["4.0": "FOURZERO", "4.1": "FOURONE",
"5.0": "FIVEZERO"]
- * - ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0":
"5.0.5"]
+ * - ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0":
"5.0.7"]
* - build-dtest-jars.sh:
* - CANDIDATE_BRANCHES=(
* "cassandra-4.0:cassandra-4.0.17"
* "cassandra-4.1:99d9faeef57c9cf5240d11eac9db5b283e45a4f9"
- * "cassandra-5.0:cassandra-5.0.5"
+ * "cassandra-5.0:cassandra-5.0.7"
*/
public enum CassandraVersion
{
@@ -185,7 +185,7 @@ public enum CassandraVersion
// NOTE: These default versions must stay in sync with
cassandraFullVersionMap in build.gradle.
String providedSupportedVersionsOrDefault =
System.getProperty("cassandra.analytics.bridges.supported_versions",
-
"cassandra-4.0.17,cassandra-5.0.5");
+
"cassandra-4.0.17,cassandra-5.0.7");
supportedVersions =
Arrays.stream(providedSupportedVersionsOrDefault.split(","))
.filter(version ->
CassandraVersion.fromVersion(version)
.filter(v
-> v.sstableFormats().contains(configuredSSTableFormat))
diff --git
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CassandraTypes.java
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CassandraTypes.java
index c31e8fe5..5444d5ed 100644
---
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CassandraTypes.java
+++
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CassandraTypes.java
@@ -37,6 +37,7 @@ import com.esotericsoftware.kryo.io.Input;
public abstract class CassandraTypes
{
public static final Pattern COLLECTION_PATTERN =
Pattern.compile("^(set|list|map|tuple)<(.+)>$", Pattern.CASE_INSENSITIVE);
+ public static final Pattern VECTOR_PATTERN =
Pattern.compile("^(vector)<(.+),(.+)>$", Pattern.CASE_INSENSITIVE);
public static final Pattern FROZEN_PATTERN =
Pattern.compile("^frozen<(.*)>$", Pattern.CASE_INSENSITIVE);
private final UDTs udts = new UDTs();
@@ -133,6 +134,8 @@ public abstract class CassandraTypes
public abstract CqlField.CqlList list(CqlField.CqlType type);
+ public abstract CqlField.CqlVector vector(CqlField.CqlType type, int
dimensions);
+
public abstract CqlField.CqlSet set(CqlField.CqlType type);
public abstract CqlField.CqlMap map(CqlField.CqlType keyType,
CqlField.CqlType valueType);
@@ -189,6 +192,14 @@ public abstract class CassandraTypes
.map(collectionType -> parseType(collectionType, udts))
.toArray(CqlField.CqlType[]::new));
}
+ Matcher vectorMatcher = VECTOR_PATTERN.matcher(type);
+ if (vectorMatcher.find())
+ {
+ // CQL vector
+ String subType = vectorMatcher.group(2);
+ int dimensions = Integer.parseInt(vectorMatcher.group(3).trim());
+ return vector(parseType(subType, udts), dimensions);
+ }
Matcher frozenMatcher = FROZEN_PATTERN.matcher(type);
if (frozenMatcher.find())
{
diff --git
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlField.java
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlField.java
index 1c15fad2..b228e36e 100644
---
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlField.java
+++
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlField.java
@@ -67,7 +67,7 @@ public class CqlField implements Serializable,
Comparable<CqlField>
{
enum InternalType
{
- NativeCql, Set, List, Map, Frozen, Udt, Tuple;
+ NativeCql, Set, List, Map, Frozen, Udt, Tuple, Vector;
public static InternalType fromString(String name)
{
@@ -77,6 +77,8 @@ public class CqlField implements Serializable,
Comparable<CqlField>
return Set;
case "list":
return List;
+ case "vector":
+ return Vector;
case "map":
return Map;
case "tuple":
@@ -237,6 +239,10 @@ public class CqlField implements Serializable,
Comparable<CqlField>
{
}
+ public interface CqlVector extends CqlCollection
+ {
+ }
+
public interface CqlTuple extends CqlCollection
{
ByteBuffer serializeTuple(Object[] values);
diff --git
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/utils/CqlUtils.java
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/utils/CqlUtils.java
index 9d989be5..60ab8c95 100644
---
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/utils/CqlUtils.java
+++
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/utils/CqlUtils.java
@@ -51,7 +51,7 @@ public final class CqlUtils
"min_index_interval",
"max_index_interval"
);
- private static final Pattern REPLICATION_FACTOR_PATTERN =
Pattern.compile("WITH REPLICATION = (\\{[^\\}]*\\})");
+ private static final Pattern REPLICATION_FACTOR_PATTERN =
Pattern.compile("WITH REPLICATION = (\\{[^\\}]*\\})", Pattern.CASE_INSENSITIVE);
// Initialize a mapper allowing single quotes to process the RF string
from the CREATE KEYSPACE statement
private static final ObjectMapper MAPPER = new
ObjectMapper().configure(JsonParser.Feature.ALLOW_SINGLE_QUOTES, true);
private static final Pattern ESCAPED_WHITESPACE_PATTERN =
Pattern.compile("(\\\\r|\\\\n|\\\\r\\n)+");
diff --git
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SqlToCqlTypeConverter.java
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SqlToCqlTypeConverter.java
index acd6721e..c67fa112 100644
---
a/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SqlToCqlTypeConverter.java
+++
b/cassandra-analytics-core/src/main/java/org/apache/cassandra/spark/bulkwriter/SqlToCqlTypeConverter.java
@@ -84,6 +84,7 @@ public final class SqlToCqlTypeConverter implements
Serializable
public static final String UDT = "udt";
public static final String VARCHAR = "varchar";
public static final String VARINT = "varint";
+ public static final String VECTOR = "vector";
private static final Logger LOGGER =
LoggerFactory.getLogger(SqlToCqlTypeConverter.class);
private static final NoOp<Object> NO_OP_CONVERTER = new NoOp<>();
private static final LongConverter LONG_CONVERTER = new LongConverter();
@@ -165,6 +166,7 @@ public final class SqlToCqlTypeConverter implements
Serializable
case TINYINT:
return NO_OP_CONVERTER;
case LIST:
+ case VECTOR:
return new ListConverter<>((CqlField.CqlCollection) cqlType);
case MAP:
assert cqlType instanceof CqlField.CqlMap;
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
index 23fd7874..5a6c4495 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
@@ -55,6 +55,7 @@ import
org.apache.cassandra.spark.transports.storage.extensions.StorageTransport
import org.apache.cassandra.spark.utils.RandomUtils;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assumptions.assumeThat;
import static org.quicktheories.QuickTheory.qt;
import static org.quicktheories.generators.SourceDSL.arbitrary;
import static org.quicktheories.generators.SourceDSL.booleans;
@@ -173,6 +174,32 @@ public class KryoSerializationTests
});
}
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testCqlFieldVector(CassandraBridge bridge)
+ {
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().withExamples(25)
+ .forAll(booleans().all(), booleans().all(),
TestUtils.cql3Type(bridge), integers().all())
+ .checkAssert((isPartitionKey, isClusteringKey, cqlType, position)
-> {
+ CqlField.CqlVector vectorType = bridge.vector(cqlType, 5);
+ CqlField field = new CqlField(isPartitionKey,
+ isClusteringKey &&
!isPartitionKey,
+ false,
+
RandomUtils.randomAlphanumeric(5, 20),
+ vectorType,
+ position);
+ Output out = serialize(bridge.getVersion(), field);
+ CqlField deserialized = deserialize(bridge.getVersion(), out,
CqlField.class);
+ assertThat(deserialized).isEqualTo(field);
+ assertThat(deserialized.name()).isEqualTo(field.name());
+ assertThat(deserialized.type()).isEqualTo(field.type());
+
assertThat(deserialized.position()).isEqualTo(field.position());
+
assertThat(deserialized.isPartitionKey()).isEqualTo(field.isPartitionKey());
+
assertThat(deserialized.isClusteringColumn()).isEqualTo(field.isClusteringColumn());
+ });
+ }
+
@ParameterizedTest
@MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
public void testCqlFieldMap(CassandraBridge bridge)
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/MockBulkWriterContext.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/MockBulkWriterContext.java
index c206bbf5..d5dab1a6 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/MockBulkWriterContext.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/MockBulkWriterContext.java
@@ -96,7 +96,7 @@ public class MockBulkWriterContext implements
BulkWriterContext, ClusterInfo, Jo
{
}
- public static final String DEFAULT_CASSANDRA_VERSION = "cassandra-5.0.5";
+ public static final String DEFAULT_CASSANDRA_VERSION = "cassandra-5.0.7";
private final UUID jobId;
private boolean skipClean = false;
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/RecordWriterTest.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/RecordWriterTest.java
index b1ce66e3..972728ec 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/RecordWriterTest.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/RecordWriterTest.java
@@ -309,7 +309,7 @@ class RecordWriterTest
@MethodSource("data")
void testWriteWithDataInMultipleSubRanges(String version)
{
- version = "cassandra-5.0.5";
+ version = "cassandra-5.0.7";
setUp(version);
MockBulkWriterContext m = Mockito.spy(writerContext);
TokenPartitioner mtp = Mockito.mock(TokenPartitioner.class);
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/StreamSessionConsistencyTest.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/StreamSessionConsistencyTest.java
index dd1fbedb..f3dcda08 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/StreamSessionConsistencyTest.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/bulkwriter/StreamSessionConsistencyTest.java
@@ -88,7 +88,7 @@ public class StreamSessionConsistencyTest
{
digestAlgorithm = new XXHash32DigestAlgorithm();
tableWriter = new MockTableWriter(folder);
- writerContext = new MockBulkWriterContext(TOKEN_RANGE_MAPPING,
"cassandra-5.0.5", consistencyLevel);
+ writerContext = new MockBulkWriterContext(TOKEN_RANGE_MAPPING,
"cassandra-5.0.7", consistencyLevel);
writerContext.setReplicationFactor(new
ReplicationFactor(NetworkTopologyStrategy, rfOptions));
transportContext = (TransportContext.DirectDataBulkWriterContext)
writerContext.transportContext();
}
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/DataTypeTests.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/DataTypeTests.java
index 210549f2..e9fa1185 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/DataTypeTests.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/DataTypeTests.java
@@ -33,17 +33,21 @@ import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
import org.apache.cassandra.bridge.CassandraBridge;
+import org.apache.cassandra.bridge.CassandraVersion;
import org.apache.cassandra.spark.TestUtils;
import org.apache.cassandra.spark.Tester;
import org.apache.cassandra.spark.data.CqlField;
import org.apache.cassandra.spark.utils.RandomUtils;
import org.apache.cassandra.spark.utils.test.TestSchema;
import org.apache.spark.sql.Row;
+import org.quicktheories.core.Gen;
import scala.collection.mutable.AbstractSeq;
import static
org.apache.cassandra.spark.utils.ScalaConversionUtils.mutableSeqAsJavaList;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assumptions.assumeThat;
import static org.quicktheories.QuickTheory.qt;
+import static org.quicktheories.generators.SourceDSL.arbitrary;
@Tag("Sequential")
public class DataTypeTests
@@ -103,6 +107,103 @@ public class DataTypeTests
);
}
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testVector(CassandraBridge bridge)
+ {
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().forAll(supportedVectorTypes(bridge))
+ .checkAssert(type ->
+ Tester.builder(TestSchema.builder(bridge)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("a",
bridge.vector(type, 10)))
+
.withExpectedRowCountPerSSTable(Tester.DEFAULT_NUM_ROWS)
+ .run(bridge.getVersion())
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testVectorVector(CassandraBridge bridge)
+ {
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().forAll(supportedVectorTypes(bridge))
+ .checkAssert(type ->
+ Tester.builder(TestSchema.builder(bridge)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("a",
bridge.vector(bridge.vector(type, 2), 5)))
+
.withExpectedRowCountPerSSTable(Tester.DEFAULT_NUM_ROWS)
+ .run(bridge.getVersion())
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testVectorList(CassandraBridge bridge)
+ {
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().forAll(supportedVectorTypes(bridge))
+ .checkAssert(type ->
+ Tester.builder(TestSchema.builder(bridge)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("a",
bridge.vector(bridge.list(type), 3)))
+
.withExpectedRowCountPerSSTable(Tester.DEFAULT_NUM_ROWS)
+ .run(bridge.getVersion())
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testVectorUDT(CassandraBridge bridge)
+ {
+ // pk -> a vector<frozen<nested_udt<x int, y type, z int>>, 10>
+ // Test vector of UDTs
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().withExamples(10)
+ .forAll(supportedVectorTypes(bridge))
+ .checkAssert(type ->
+ Tester.builder(TestSchema.builder(bridge)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("a",
bridge.vector(
+
bridge.udt("keyspace", "nested_udt")
+
.withField("x", bridge.aInt())
+
.withField("y", type)
+
.withField("z", bridge.aInt())
+
.build().frozen(),
+
10)))
+ .run(bridge.getVersion())
+ );
+ }
+
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testVectorTuple(CassandraBridge bridge)
+ {
+ // pk -> a vector<frozen<tuple<type, float, text>>, 7>
+ // Test tuple nested within vector
+
assumeThat(bridge.getVersion().versionNumber()).isGreaterThanOrEqualTo(CassandraVersion.FIVEZERO.versionNumber());
+ qt().withExamples(10)
+ .forAll(supportedVectorTypes(bridge))
+ .checkAssert(type ->
+ Tester.builder(TestSchema.builder(bridge)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("a",
bridge.vector(bridge.tuple(type,
+
bridge.aFloat(),
+
bridge.text()).frozen(), 7)))
+ .run(bridge.getVersion())
+ );
+ }
+
+ private static Gen<CqlField.NativeType>
supportedVectorTypes(CassandraBridge bridge)
+ {
+ // TODO: Vector of list of durations fail, because we cannot replace
DurationSerializer with
+ // AnalyticsDurationSerializer across all serializers used by
VectorType.
+ List<CqlField.NativeType> supportedTypes =
bridge.supportedTypes().stream()
+ .filter(t ->
!t.equals(bridge.duration()))
+
.collect(Collectors.toList());
+ return arbitrary().pick(supportedTypes);
+ }
+
@ParameterizedTest
@MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
public void testList(CassandraBridge bridge)
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/MiscTests.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/MiscTests.java
index ac7a1f18..b8259c29 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/MiscTests.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/endtoend/MiscTests.java
@@ -158,7 +158,7 @@ public class MiscTests
public void testQuotedKeyspaceName(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withPartitionKey("pk",
bridge.uuid())
.withColumn("c1",
bridge.varint())
.withColumn("c2", bridge.text())
@@ -184,7 +184,7 @@ public class MiscTests
public void testQuotedTableName(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("pk",
bridge.uuid())
.withColumn("c1",
bridge.varint())
@@ -198,7 +198,7 @@ public class MiscTests
public void testReservedWordTableName(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("table")
.withPartitionKey("pk",
bridge.uuid())
.withColumn("c1",
bridge.varint())
@@ -212,7 +212,7 @@ public class MiscTests
public void testQuotedPartitionKey(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("Partition_Key_0", bridge.uuid())
.withColumn("c1",
bridge.varint())
@@ -226,7 +226,7 @@ public class MiscTests
public void testMultipleQuotedPartitionKeys(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("Partition_Key_0", bridge.uuid())
.withPartitionKey("Partition_Key_1", bridge.bigint())
@@ -243,7 +243,7 @@ public class MiscTests
public void testQuotedPartitionClusteringKeys(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("a",
bridge.uuid())
.withClusteringKey("Clustering_Key_0", bridge.bigint())
@@ -258,7 +258,7 @@ public class MiscTests
public void testQuotedColumnNames(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("Partition_Key_0", bridge.uuid())
.withColumn("Column_1",
bridge.varint())
@@ -272,7 +272,7 @@ public class MiscTests
public void testQuotedColumnNamesWithColumnFilter(CassandraBridge bridge)
{
Tester.builder(keyspace1 -> TestSchema.builder(bridge)
- .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", "_"))
+ .withKeyspace("Quoted_Keyspace_"
+ UUID.randomUUID().toString().replaceAll("-", ""))
.withTable("Quoted_Table_" +
UUID.randomUUID().toString().replaceAll("-", "_"))
.withPartitionKey("Partition_Key_0", bridge.uuid())
.withColumn("Column_1",
bridge.varint())
diff --git a/cassandra-analytics-integration-framework/build.gradle
b/cassandra-analytics-integration-framework/build.gradle
index 6c1c65b0..7c97cb3b 100644
--- a/cassandra-analytics-integration-framework/build.gradle
+++ b/cassandra-analytics-integration-framework/build.gradle
@@ -32,7 +32,7 @@ if (propertyWithDefault("artifactType", null) == "spark")
apply from: "$rootDir/gradle/common/publishing.gradle"
}
-ext.dtestJar = System.getenv("DTEST_JAR") ?: "dtest-5.0.5.jar" // latest
supported Cassandra build is 5.0
+ext.dtestJar = System.getenv("DTEST_JAR") ?: "dtest-5.0.7.jar" // latest
supported Cassandra build is 5.0
def dtestJarFullPath = "${dependencyLocation}${ext.dtestJar}"
test {
@@ -50,7 +50,7 @@ dependencies {
// classpath while running integration tests. Instead, a dedicated
classloader will load the
// dtest jar while provisioning the in-jvm dtest Cassandra cluster
compileOnly(files("${dtestJarFullPath}"))
- api("org.apache.cassandra:dtest-api:0.0.16")
+ api("org.apache.cassandra:dtest-api:0.0.18")
// Needed by the Cassandra dtest framework
// JUnit
api("org.junit.jupiter:junit-jupiter-api:${project.junitVersion}")
diff --git
a/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/distributed/impl/CassandraCluster.java
b/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/distributed/impl/CassandraCluster.java
index 93a877ea..69686b90 100644
---
a/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/distributed/impl/CassandraCluster.java
+++
b/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/distributed/impl/CassandraCluster.java
@@ -213,6 +213,11 @@ public class CassandraCluster<I extends IInstance>
implements IClusterExtension<
return delegate.newInstanceConfig();
}
+ public IInstanceConfig createInstanceConfig(int i)
+ {
+ throw new UnsupportedOperationException();
+ }
+
@Override
public ICluster<I> delegate()
{
diff --git
a/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/SharedClusterIntegrationTestBase.java
b/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/SharedClusterIntegrationTestBase.java
index a1004686..0577e429 100644
---
a/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/SharedClusterIntegrationTestBase.java
+++
b/cassandra-analytics-integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/SharedClusterIntegrationTestBase.java
@@ -503,6 +503,17 @@ public abstract class SharedClusterIntegrationTestBase
return queryAllDataWithDriver(table, ConsistencyLevel.ALL);
}
+ /**
+ * Convenience method to count rows from the provided {@code table} at
consistency level ALL.
+ *
+ * @param table the qualified Cassandra table name
+ * @return all the data queried from the table
+ */
+ protected Long countDataWithDriver(QualifiedName table)
+ {
+ return countDataWithDriver(table, ConsistencyLevel.ALL);
+ }
+
/**
* Convenience method to query all data from the provided {@code table} at
the specified consistency level.
*
@@ -512,11 +523,31 @@ public abstract class SharedClusterIntegrationTestBase
*/
protected ResultSet queryAllDataWithDriver(QualifiedName table,
ConsistencyLevel consistency)
{
- Cluster driverCluster = createDriverCluster(cluster.delegate());
- Session session = driverCluster.connect();
- SimpleStatement statement = new SimpleStatement(String.format("SELECT
* FROM %s;", table));
-
statement.setConsistencyLevel(com.datastax.driver.core.ConsistencyLevel.valueOf(consistency.name()));
- return session.execute(statement);
+ try (Cluster driverCluster = createDriverCluster(cluster.delegate());
+ Session session = driverCluster.connect())
+ {
+ SimpleStatement statement = new
SimpleStatement(String.format("SELECT * FROM %s;", table));
+
statement.setConsistencyLevel(com.datastax.driver.core.ConsistencyLevel.valueOf(consistency.name()));
+ return session.execute(statement);
+ }
+ }
+
+ /**
+ * Convenience method to count rows from the provided {@code table} at the
specified consistency level.
+ *
+ * @param table the qualified Cassandra table name
+ * @param consistency the consistency level to use for querying the data
+ * @return record count
+ */
+ protected Long countDataWithDriver(QualifiedName table, ConsistencyLevel
consistency)
+ {
+ try (Cluster driverCluster = createDriverCluster(cluster.delegate());
+ Session session = driverCluster.connect())
+ {
+ SimpleStatement statement = new
SimpleStatement(String.format("SELECT COUNT(*) FROM %s;", table));
+
statement.setConsistencyLevel(com.datastax.driver.core.ConsistencyLevel.valueOf(consistency.name()));
+ return session.execute(statement).one().getLong(0);
+ }
}
// Utility methods
diff --git
a/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkReaderVectorTest.java
b/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkReaderVectorTest.java
new file mode 100644
index 00000000..14c93ed3
--- /dev/null
+++
b/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkReaderVectorTest.java
@@ -0,0 +1,115 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.cassandra.analytics;
+
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.List;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.stream.Collectors;
+
+import org.junit.jupiter.api.Test;
+
+import com.vdurmont.semver4j.Semver;
+import org.apache.cassandra.distributed.api.ConsistencyLevel;
+import org.apache.cassandra.distributed.api.IInstance;
+import org.apache.cassandra.sidecar.testing.QualifiedName;
+import org.apache.cassandra.testing.ClusterBuilderConfiguration;
+import org.apache.cassandra.testing.TestUtils;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Row;
+
+import static org.apache.cassandra.testing.TestUtils.DC1_RF1;
+import static org.apache.cassandra.testing.TestUtils.TEST_KEYSPACE;
+import static org.apache.cassandra.testing.TestUtils.uniqueTestTableFullName;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assumptions.assumeThat;
+
+/**
+ * Tests bulk reader functionality
+ */
+class BulkReaderVectorTest extends SharedClusterSparkIntegrationTestBase
+{
+ static final int ROW_COUNT = 10;
+ static final int DIMENSIONS = 3;
+ static final List<List<Float>> DATASET = new ArrayList<>();
+ static QualifiedName table1 = uniqueTestTableFullName(TEST_KEYSPACE);
+
+ static
+ {
+ for (int i = 0; i < ROW_COUNT; i++)
+ {
+ List<Float> vector = new ArrayList<>();
+ for (int j = 0; j < DIMENSIONS; j++)
+ {
+ vector.add(ThreadLocalRandom.current().nextFloat());
+ }
+ DATASET.add(vector);
+ }
+ }
+
+ @Override
+ protected void beforeClusterProvisioning()
+ {
+
assumeThat(TestUtils.getDTestClusterVersion().isGreaterThanOrEqualTo(new
Semver("5.0", Semver.SemverType.LOOSE)))
+ .describedAs("Vector type was introduced in Cassandra 5.0")
+ .isTrue();
+ }
+
+ @Override
+ protected ClusterBuilderConfiguration testClusterConfiguration()
+ {
+ return super.testClusterConfiguration()
+ .nodesPerDc(2);
+ }
+
+ @Test
+ void testReadingVectorColumn()
+ {
+ Dataset<Row> data = bulkReaderDataFrame(table1).load();
+
+ List<Row> rows = data.collectAsList().stream()
+ .sorted(Comparator.comparing(row ->
row.getInt(0)))
+ .collect(Collectors.toList());
+ assertThat(rows.size()).isEqualTo(ROW_COUNT);
+
+ for (int i = 0; i < ROW_COUNT; i++)
+ {
+ Row row = rows.get(i);
+ List<Float> value = DATASET.get(i);
+ assertThat(row.getList(1)).isEqualTo(value);
+ }
+ }
+
+ @Override
+ protected void initializeSchemaForTest()
+ {
+ createTestKeyspace(TEST_KEYSPACE, DC1_RF1);
+ createTestTable(table1, "CREATE TABLE IF NOT EXISTS %s (id int PRIMARY
KEY, value vector<float, " + DIMENSIONS + ">);");
+
+ IInstance firstRunningInstance = cluster.getFirstRunningInstance();
+ for (int i = 0; i < ROW_COUNT; i++)
+ {
+ List<Float> value = DATASET.get(i);
+ String query = String.format("INSERT INTO %s (id, value) VALUES
(%d, %s);", table1, i, value.toString());
+ firstRunningInstance.coordinator().execute(query,
ConsistencyLevel.ALL);
+ }
+ }
+}
diff --git
a/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkWriteVectorTest.java
b/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkWriteVectorTest.java
new file mode 100644
index 00000000..3571b159
--- /dev/null
+++
b/cassandra-analytics-integration-tests/src/test/java/org/apache/cassandra/analytics/BulkWriteVectorTest.java
@@ -0,0 +1,115 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.cassandra.analytics;
+
+import org.junit.jupiter.api.Test;
+
+import com.vdurmont.semver4j.Semver;
+import org.apache.cassandra.distributed.api.ConsistencyLevel;
+import org.apache.cassandra.distributed.api.ICoordinator;
+import org.apache.cassandra.sidecar.testing.QualifiedName;
+import org.apache.cassandra.testing.ClusterBuilderConfiguration;
+import org.apache.cassandra.testing.TestUtils;
+import org.apache.spark.sql.Dataset;
+import org.apache.spark.sql.Row;
+
+import static org.apache.cassandra.testing.TestUtils.DC1_RF3;
+import static org.apache.cassandra.testing.TestUtils.ROW_COUNT;
+import static org.apache.cassandra.testing.TestUtils.TEST_KEYSPACE;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assumptions.assumeThat;
+
+public class BulkWriteVectorTest extends SharedClusterSparkIntegrationTestBase
+{
+ static final QualifiedName VECTOR_TABLE_NAME = new
QualifiedName(TEST_KEYSPACE, "test_vector");
+ public static final String VECTOR_TABLE_CREATE = "CREATE TABLE " +
VECTOR_TABLE_NAME + " (\n"
+ + " id BIGINT
PRIMARY KEY,\n"
+ + " value
vector<FLOAT, 3>);";
+
+ private ICoordinator coordinator;
+
+ @Test
+ void testVectorOfFloats()
+ {
+ int numRowsInserted = populateVectorOfFloats();
+ // Create a spark frame with the data inserted during the setup
+ Dataset<Row> sourceData =
bulkReaderDataFrame(VECTOR_TABLE_NAME).load();
+ assertThat(sourceData.count()).isEqualTo(numRowsInserted);
+
+ // truncate table to re-insert the data
+ truncateTable(VECTOR_TABLE_NAME);
+
+ // Insert the dataset containing vectors
+ bulkWriterDataFrameWriter(sourceData, VECTOR_TABLE_NAME).save();
+
+ // Count rows because Java driver 3.x cannot read vector type
+
assertThat(countDataWithDriver(VECTOR_TABLE_NAME)).isEqualTo(numRowsInserted);
+ }
+
+ private int populateVectorOfFloats()
+ {
+ String insert = "INSERT INTO %s (id, value) VALUES (%d, [%f, %f, %f])";
+
+ int i = 0;
+ for (; i < ROW_COUNT; i++)
+ {
+ float j = (float) i;
+ cluster.schemaChangeIgnoringStoppedInstances(String.format(insert,
VECTOR_TABLE_NAME,
+ i, j,
j, j));
+ }
+
+ // test null value
+ coordinator.execute(String.format("insert into %s (id) values (%d)",
+ VECTOR_TABLE_NAME, i++),
ConsistencyLevel.ALL);
+
+ return i;
+ }
+
+ @Override
+ protected ClusterBuilderConfiguration testClusterConfiguration()
+ {
+ return super.testClusterConfiguration()
+ .nodesPerDc(3);
+ }
+
+ @Override
+ protected void beforeClusterProvisioning()
+ {
+
assumeThat(TestUtils.getDTestClusterVersion().isGreaterThanOrEqualTo(new
Semver("5.0", Semver.SemverType.LOOSE)))
+ .describedAs("Vector type was introduced in Cassandra 5.0")
+ .isTrue();
+ }
+
+ @Override
+ protected void initializeSchemaForTest()
+ {
+ coordinator = cluster.getFirstRunningInstance().coordinator();
+
+ createTestKeyspace(VECTOR_TABLE_NAME, DC1_RF3);
+
+ cluster.schemaChangeIgnoringStoppedInstances(VECTOR_TABLE_CREATE);
+ }
+
+ private void truncateTable(QualifiedName tableName)
+ {
+ cluster.schemaChangeIgnoringStoppedInstances(String.format(
+ "TRUNCATE %s.%s",
+ TEST_KEYSPACE, tableName.table()));
+ }
+}
diff --git
a/cassandra-analytics-spark-four-zero-converter/src/main/java/org/apache/cassandra/spark/data/converter/SparkSqlTypeConverterImplementation.java
b/cassandra-analytics-spark-four-zero-converter/src/main/java/org/apache/cassandra/spark/data/converter/SparkSqlTypeConverterImplementation.java
index 5c4b9352..6610a2c3 100644
---
a/cassandra-analytics-spark-four-zero-converter/src/main/java/org/apache/cassandra/spark/data/converter/SparkSqlTypeConverterImplementation.java
+++
b/cassandra-analytics-spark-four-zero-converter/src/main/java/org/apache/cassandra/spark/data/converter/SparkSqlTypeConverterImplementation.java
@@ -166,6 +166,10 @@ public class SparkSqlTypeConverterImplementation
implements SparkSqlTypeConverte
{
return new SparkSet(INSTANCE, (CqlField.CqlSet) cqlType);
}
+ else if (cqlType instanceof CqlField.CqlVector)
+ {
+ return new SparkList(INSTANCE, (CqlField.CqlVector) cqlType);
+ }
else if (cqlType instanceof CqlField.CqlList)
{
return new SparkList(INSTANCE, (CqlField.CqlList) cqlType);
diff --git
a/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CassandraBridge.java
b/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CassandraBridge.java
index 67c5bc7a..22db7f2f 100644
---
a/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CassandraBridge.java
+++
b/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CassandraBridge.java
@@ -318,6 +318,11 @@ public abstract class CassandraBridge
return cassandraTypes().list(type);
}
+ public CqlField.CqlVector vector(CqlField.CqlType type, int dimensions)
+ {
+ return cassandraTypes().vector(type, dimensions);
+ }
+
public CqlField.CqlSet set(CqlField.CqlType type)
{
return cassandraTypes().set(type);
diff --git
a/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/spark/data/converter/types/VectorTypeTests.java
b/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/spark/data/converter/types/VectorTypeTests.java
new file mode 100644
index 00000000..fc44485c
--- /dev/null
+++
b/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/spark/data/converter/types/VectorTypeTests.java
@@ -0,0 +1,59 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.cassandra.spark.data.converter.types;
+
+import java.util.List;
+import java.util.Set;
+
+import org.junit.jupiter.api.Test;
+
+import org.apache.cassandra.bridge.CassandraBridgeImplementation;
+import org.apache.cassandra.spark.data.complex.CqlList;
+import org.apache.cassandra.spark.data.complex.CqlVector;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+public class VectorTypeTests
+{
+ private static final CassandraBridgeImplementation BRIDGE = new
CassandraBridgeImplementation();
+
+ @Test
+ public void testSimpleTypeConversion()
+ {
+ CqlVector cqlVector = new
CqlVector(org.apache.cassandra.spark.data.types.Float.INSTANCE, 3);
+ Object cqlWriterObj = cqlVector.convertForCqlWriter(List.of(3.14f,
0.0f, -1f), BRIDGE.getVersion(), false);
+ assertThat(cqlWriterObj).isInstanceOf(List.class);
+ List<Float> cqlWriterList = (List<Float>) cqlWriterObj;
+ assertThat(cqlWriterList).containsExactly(3.14f, 0.0f, -1f);
+ }
+
+ @Test
+ public void testComplexTypeConversion()
+ {
+ CqlVector cqlVector = new
CqlVector(CqlList.set(org.apache.cassandra.spark.data.types.Float.INSTANCE), 3);
+ Object cqlWriterObj =
cqlVector.convertForCqlWriter(List.of(Set.of(3.14f, 0f), Set.of(1f), Set.of()),
BRIDGE.getVersion(), false);
+ assertThat(cqlWriterObj).isInstanceOf(List.class);
+ List<? extends Set<Float>> cqlWriterList = (List<? extends
Set<Float>>) cqlWriterObj;
+ assertThat(cqlWriterList).hasSize(3);
+ assertThat(cqlWriterList.get(0)).containsExactlyInAnyOrder(3.14f, 0f);
+ assertThat(cqlWriterList.get(1)).containsExactly(1f);
+ assertThat(cqlWriterList.get(2)).isEmpty();
+ }
+}
diff --git
a/cassandra-five-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraTypesImplementation.java
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraTypesImplementation.java
index 6d5804cb..085afed1 100644
---
a/cassandra-five-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraTypesImplementation.java
+++
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraTypesImplementation.java
@@ -24,6 +24,7 @@ import java.nio.file.Files;
import java.nio.file.Path;
import java.util.UUID;
+import com.esotericsoftware.kryo.io.Input;
import org.apache.cassandra.config.Config;
import org.apache.cassandra.config.DataStorageSpec;
import org.apache.cassandra.config.DatabaseDescriptor;
@@ -32,6 +33,8 @@ import
org.apache.cassandra.db.commitlog.CommitLogSegmentManagerStandard;
import org.apache.cassandra.dht.Murmur3Partitioner;
import org.apache.cassandra.locator.SimpleSnitch;
import org.apache.cassandra.security.EncryptionContext;
+import org.apache.cassandra.spark.data.CqlField;
+import org.apache.cassandra.spark.data.complex.CqlVector;
public class CassandraTypesImplementation extends AbstractCassandraTypes
{
@@ -88,4 +91,20 @@ public class CassandraTypesImplementation extends
AbstractCassandraTypes
DatabaseDescriptor.getRawConfig().commitlog_total_space = new
DataStorageSpec.IntMebibytesBound(1024);
DatabaseDescriptor.setCommitLogSegmentMgrProvider(commitLog -> new
CommitLogSegmentManagerStandard(commitLog, commitLogPath.toString()));
}
+
+ @Override
+ public CqlField.CqlType readType(CqlField.CqlType.InternalType type, Input
input)
+ {
+ if (type == CqlField.CqlType.InternalType.Vector)
+ {
+ return CqlVector.read(input, this);
+ }
+ return super.readType(type, input);
+ }
+
+ @Override
+ public CqlField.CqlVector vector(CqlField.CqlType type, int dimensions)
+ {
+ return new CqlVector(type, dimensions);
+ }
}
diff --git
a/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlVector.java
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlVector.java
new file mode 100644
index 00000000..cef8a1b8
--- /dev/null
+++
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlVector.java
@@ -0,0 +1,166 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.cassandra.spark.data.complex;
+
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+import com.google.common.base.Preconditions;
+
+import com.esotericsoftware.kryo.io.Input;
+import com.esotericsoftware.kryo.io.Output;
+import org.apache.cassandra.bridge.CassandraVersion;
+import org.apache.cassandra.cql3.functions.types.SettableByIndexData;
+import org.apache.cassandra.db.marshal.AbstractType;
+import org.apache.cassandra.db.marshal.VectorType;
+import org.apache.cassandra.db.rows.CellPath;
+import org.apache.cassandra.serializers.TypeSerializer;
+import org.apache.cassandra.spark.data.CassandraTypes;
+import org.apache.cassandra.spark.data.CqlField;
+import org.apache.cassandra.spark.data.CqlType;
+import org.apache.cassandra.utils.TimeUUID;
+import org.jetbrains.annotations.NotNull;
+
+public class CqlVector extends CqlCollection implements CqlField.CqlVector
+{
+ private final int dimensions;
+
+ public CqlVector(CqlField.CqlType type, int dimensions)
+ {
+ super(type);
+ this.dimensions = dimensions;
+ this.hashCode = Objects.hash(this.hashCode, dimensions);
+ }
+
+ public static CqlVector read(Input input, CassandraTypes cassandraTypes)
+ {
+ int dimensions = input.readInt();
+ CqlField.CqlType[] types = CqlCollection.readTypes(input,
cassandraTypes);
+ Preconditions.checkArgument(types.length == 1, "Unexpected number of
vector subtypes: " + types.length);
+ return new CqlVector(types[0], dimensions);
+ }
+
+ @Override
+ public void write(Output output)
+ {
+ CqlField.CqlType.write(this, output);
+ output.writeInt(dimensions);
+ writeTypes(output);
+ }
+
+ @Override
+ public AbstractType<?> dataType(boolean isMultiCell)
+ {
+ return VectorType.getInstance(((CqlType) type()).dataType(),
dimensions);
+ }
+
+ @Override
+ public InternalType internalType()
+ {
+ return InternalType.Vector;
+ }
+
+ @Override
+ @SuppressWarnings("unchecked")
+ public <T> TypeSerializer<T> serializer()
+ {
+ return (TypeSerializer<T>) dataType(false).getSerializer();
+ }
+
+ @Override
+ public String name()
+ {
+ return "vector";
+ }
+
+ @Override
+ public String cqlName()
+ {
+ return String.format("%s<%s, %d>",
+ internalType().name().toLowerCase(),
+ types.get(0).cqlName(),
+ dimensions);
+ }
+
+ @Override
+ protected void setInnerValueInternal(SettableByIndexData<?> udtValue, int
position, @NotNull Object value)
+ {
+ List<?> vector = (List<?>) value;
+ validate(vector);
+ udtValue.setVector(position, vector);
+ }
+
+ @Override
+ public Object randomValue(int minCollectionSize)
+ {
+ return IntStream.range(0, dimensions)
+ .mapToObj(element ->
type().randomValue(minCollectionSize))
+ .collect(Collectors.toList());
+ }
+
+ @Override
+ public org.apache.cassandra.cql3.functions.types.DataType
driverDataType(boolean isFrozen)
+ {
+ return
org.apache.cassandra.cql3.functions.types.DataType.vector(((CqlType)
type()).driverDataType(isFrozen), dimensions);
+ }
+
+ @Override
+ public Object convertForCqlWriter(Object value, CassandraVersion version,
boolean isCollectionElement)
+ {
+ List<?> vector = (List<?>) value;
+ validate(vector);
+ return vector.stream()
+ .map(element -> type().convertForCqlWriter(element,
version, true))
+ .collect(Collectors.toList());
+ }
+
+ @Override
+ public int hashCode()
+ {
+ return super.hashCode();
+ }
+
+ @Override
+ public boolean equals(Object o)
+ {
+ if (this == o)
+ {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass())
+ {
+ return false;
+ }
+ CqlVector that = (CqlVector) o;
+ return super.equals(o) && dimensions == that.dimensions;
+ }
+
+ protected CellPath randomCellPath()
+ {
+ return CellPath.create(TimeUUID.Generator.nextTimeUUID().toBytes());
+ }
+
+ private void validate(List<?> vector)
+ {
+ Preconditions.checkArgument(vector.size() == dimensions, "Expected " +
dimensions + " for vector: " + vector);
+ }
+}
diff --git
a/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
new file mode 100644
index 00000000..330331f0
--- /dev/null
+++
b/cassandra-five-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
@@ -0,0 +1,114 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.cassandra.spark.reader;
+
+import java.util.Collections;
+import java.util.Set;
+import java.util.UUID;
+import java.util.function.Function;
+
+import com.google.common.annotations.VisibleForTesting;
+
+import org.apache.cassandra.cql3.CQL3Type;
+import org.apache.cassandra.db.marshal.AbstractType;
+import org.apache.cassandra.db.marshal.VectorType;
+import org.apache.cassandra.spark.data.CassandraTypes;
+import org.apache.cassandra.spark.data.CqlTable;
+import org.apache.cassandra.spark.data.ReplicationFactor;
+import org.apache.cassandra.spark.data.partitioner.Partitioner;
+import org.jetbrains.annotations.Nullable;
+
+public class SchemaBuilder extends AbstractSchemaBuilder
+{
+ public SchemaBuilder(CqlTable table, Partitioner partitioner, boolean
enableCdc)
+ {
+ this(table, partitioner, null, enableCdc);
+ }
+
+ public SchemaBuilder(CqlTable table, Partitioner partitioner)
+ {
+ this(table, partitioner, null, false);
+ }
+
+ public SchemaBuilder(CqlTable table, Partitioner partitioner, UUID
tableId, boolean enableCdc)
+ {
+ this(table.createStatement(),
+ table.keyspace(),
+ table.replicationFactor(),
+ partitioner,
+ table::udtCreateStmts,
+ tableId,
+ 0,
+ enableCdc);
+ }
+
+ @VisibleForTesting
+ public SchemaBuilder(String createStmt, String keyspace, ReplicationFactor
replicationFactor)
+ {
+ this(createStmt, keyspace, replicationFactor,
Partitioner.Murmur3Partitioner, bridge -> Collections.emptySet(), null, 0,
false);
+ }
+
+ @VisibleForTesting
+ public SchemaBuilder(String createStmt,
+ String keyspace,
+ ReplicationFactor replicationFactor,
+ Partitioner partitioner)
+ {
+ this(createStmt, keyspace, replicationFactor, partitioner, bridge ->
Collections.emptySet(), null, 0, false);
+ }
+
+ public SchemaBuilder(String createStmt,
+ String keyspace,
+ ReplicationFactor replicationFactor,
+ Partitioner partitioner,
+ Function<CassandraTypes, Set<String>>
udtStatementsProvider,
+ @Nullable UUID tableId,
+ int indexCount,
+ boolean enableCdc)
+ {
+ super(createStmt, keyspace, replicationFactor, partitioner,
udtStatementsProvider,
+ tableId, indexCount, enableCdc);
+ }
+
+ @Override
+ protected void validateType(CQL3Type cqlType)
+ {
+ if (!(cqlType instanceof CQL3Type.Native)
+ && !(cqlType instanceof CQL3Type.Collection)
+ && !(cqlType instanceof CQL3Type.UserDefined)
+ && !(cqlType instanceof CQL3Type.Tuple)
+ && !(cqlType instanceof CQL3Type.Vector))
+ {
+ throw new UnsupportedOperationException("Only native, collection,
tuples, vectors or UDT data types are supported, "
+ + "unsupported data type:
" + cqlType.toString());
+ }
+ if (cqlType instanceof CQL3Type.Vector)
+ {
+ CQL3Type.Vector vector = (CQL3Type.Vector) cqlType;
+ VectorType<?> vectorType = vector.getType();
+ for (AbstractType<?> subType : vectorType.subTypes())
+ {
+ validateType(subType);
+ }
+ return;
+ }
+ super.validateType(cqlType);
+ }
+}
diff --git
a/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/spark/reader/AbstractStreamScanner.java
b/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/spark/reader/AbstractStreamScanner.java
index 0853325c..2e7d30d0 100644
---
a/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/spark/reader/AbstractStreamScanner.java
+++
b/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/spark/reader/AbstractStreamScanner.java
@@ -384,7 +384,7 @@ public abstract class AbstractStreamScanner implements
StreamScanner<RowData>, C
{
boolean isStatic = cell.column().isStatic();
rowData.setColumnNameCopy(ReaderUtils.encodeCellName(metadata,
- isStatic ?
Clustering.STATIC_CLUSTERING : clustering,
+ isStatic ?
Clustering.STATIC_CLUSTERING : clustering,
cell.column().name.bytes,
null));
if (cell.isTombstone())
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/AbstractCassandraTypes.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/AbstractCassandraTypes.java
index 8c5537e0..6dfeea89 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/AbstractCassandraTypes.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/AbstractCassandraTypes.java
@@ -253,6 +253,12 @@ public abstract class AbstractCassandraTypes extends
CassandraTypes
return CqlCollection.list(type);
}
+ @Override
+ public CqlField.CqlVector vector(CqlField.CqlType type, int dimensions)
+ {
+ throw new UnsupportedOperationException("Vector data type is available
in C* 5.x.");
+ }
+
@Override
public CqlField.CqlSet set(CqlField.CqlType type)
{
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlCollection.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlCollection.java
index 4f7f9b45..34440a83 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlCollection.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/data/complex/CqlCollection.java
@@ -39,7 +39,7 @@ import org.apache.cassandra.spark.data.CqlType;
public abstract class CqlCollection extends CqlType implements
CqlField.CqlCollection
{
public final List<CqlField.CqlType> types;
- private final int hashCode;
+ protected int hashCode;
CqlCollection(CqlField.CqlType type)
{
@@ -174,6 +174,12 @@ public abstract class CqlCollection extends CqlType
implements CqlField.CqlColle
}
public static CqlCollection read(CqlField.CqlType.InternalType
internalType, Input input, CassandraTypes cassandraTypes)
+ {
+ CqlField.CqlType[] types = readTypes(input, cassandraTypes);
+ return CqlCollection.build(internalType, types);
+ }
+
+ protected static CqlField.CqlType[] readTypes(Input input, CassandraTypes
cassandraTypes)
{
int numTypes = input.readInt();
CqlField.CqlType[] types = new CqlField.CqlType[numTypes];
@@ -181,13 +187,18 @@ public abstract class CqlCollection extends CqlType
implements CqlField.CqlColle
{
types[type] = CqlField.CqlType.read(input, cassandraTypes);
}
- return CqlCollection.build(internalType, types);
+ return types;
}
@Override
public void write(Output output)
{
CqlField.CqlType.write(this, output);
+ writeTypes(output);
+ }
+
+ protected void writeTypes(Output output)
+ {
output.writeInt(this.types.size());
for (CqlField.CqlType type : this.types)
{
@@ -212,7 +223,7 @@ public abstract class CqlCollection extends CqlType
implements CqlField.CqlColle
{
return true;
}
- if (this.getClass() != other.getClass())
+ if (!(other instanceof CqlCollection))
{
return false;
}
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/AbstractSchemaBuilder.java
similarity index 92%
copy from
cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
copy to
cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/AbstractSchemaBuilder.java
index b565565b..e353ceaf 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/AbstractSchemaBuilder.java
@@ -74,30 +74,30 @@ import org.apache.cassandra.utils.Pair;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
-public class SchemaBuilder
+public abstract class AbstractSchemaBuilder
{
- private static final Logger LOGGER =
LoggerFactory.getLogger(SchemaBuilder.class);
-
- private final TableMetadata metadata;
- private final KeyspaceMetadata keyspaceMetadata;
- private final String createStmt;
- private final String keyspace;
- private final ReplicationFactor replicationFactor;
- private final CassandraTypes cassandraTypes;
- private final int indexCount;
- private final boolean enableCdc;
-
- public SchemaBuilder(CqlTable table, Partitioner partitioner, boolean
enableCdc)
+ private static final Logger LOGGER =
LoggerFactory.getLogger(AbstractSchemaBuilder.class);
+
+ protected final TableMetadata metadata;
+ protected final KeyspaceMetadata keyspaceMetadata;
+ protected final String createStmt;
+ protected final String keyspace;
+ protected final ReplicationFactor replicationFactor;
+ protected final CassandraTypes cassandraTypes;
+ protected final int indexCount;
+ protected final boolean enableCdc;
+
+ public AbstractSchemaBuilder(CqlTable table, Partitioner partitioner,
boolean enableCdc)
{
this(table, partitioner, null, enableCdc);
}
- public SchemaBuilder(CqlTable table, Partitioner partitioner)
+ public AbstractSchemaBuilder(CqlTable table, Partitioner partitioner)
{
this(table, partitioner, null, table.cdc());
}
- public SchemaBuilder(CqlTable table, Partitioner partitioner, UUID
tableId, boolean enableCdc)
+ public AbstractSchemaBuilder(CqlTable table, Partitioner partitioner, UUID
tableId, boolean enableCdc)
{
this(table.createStatement(),
table.keyspace(),
@@ -110,28 +110,28 @@ public class SchemaBuilder
}
@VisibleForTesting
- public SchemaBuilder(String createStmt, String keyspace, ReplicationFactor
replicationFactor)
+ public AbstractSchemaBuilder(String createStmt, String keyspace,
ReplicationFactor replicationFactor)
{
this(createStmt, keyspace, replicationFactor,
Partitioner.Murmur3Partitioner, bridge -> Collections.emptySet(), null, 0,
false);
}
@VisibleForTesting
- public SchemaBuilder(String createStmt,
- String keyspace,
- ReplicationFactor replicationFactor,
- Partitioner partitioner)
+ public AbstractSchemaBuilder(String createStmt,
+ String keyspace,
+ ReplicationFactor replicationFactor,
+ Partitioner partitioner)
{
this(createStmt, keyspace, replicationFactor, partitioner, bridge ->
Collections.emptySet(), null, 0, false);
}
- public SchemaBuilder(String createStmt,
- String keyspace,
- ReplicationFactor replicationFactor,
- Partitioner partitioner,
- Function<CassandraTypes, Set<String>>
udtStatementsProvider,
- @Nullable UUID tableId,
- int indexCount,
- boolean enableCdc)
+ public AbstractSchemaBuilder(String createStmt,
+ String keyspace,
+ ReplicationFactor replicationFactor,
+ Partitioner partitioner,
+ Function<CassandraTypes, Set<String>>
udtStatementsProvider,
+ @Nullable UUID tableId,
+ int indexCount,
+ boolean enableCdc)
{
this.createStmt = createStmt;
this.keyspace = keyspace;
@@ -234,20 +234,20 @@ public class SchemaBuilder
validateType(column.type);
}
- private void validateType(AbstractType<?> type)
+ protected void validateType(AbstractType<?> type)
{
validateType(type.asCQL3Type());
}
- private void validateType(CQL3Type cqlType)
+ protected void validateType(CQL3Type cqlType)
{
if (!(cqlType instanceof CQL3Type.Native)
- && !(cqlType instanceof CQL3Type.Collection)
- && !(cqlType instanceof CQL3Type.UserDefined)
- && !(cqlType instanceof CQL3Type.Tuple))
+ && !(cqlType instanceof CQL3Type.Collection)
+ && !(cqlType instanceof CQL3Type.UserDefined)
+ && !(cqlType instanceof CQL3Type.Tuple))
{
throw new UnsupportedOperationException("Only native, collection,
tuples or UDT data types are supported, "
- + "unsupported data type: "
+ cqlType.toString());
+ + "unsupported data type:
" + cqlType.toString());
}
if (cqlType instanceof CQL3Type.Native)
@@ -479,8 +479,7 @@ public class SchemaBuilder
replicationFactor,
fields,
new HashSet<>(udts.values()),
- indexCount,
- enableCdc);
+ indexCount);
}
private Map<String, CqlField.CqlUdt> buildsUdts(KeyspaceMetadata
keyspaceMetadata)
@@ -491,7 +490,7 @@ public class SchemaBuilder
while (!userTypes.isEmpty())
{
UserType userType = userTypes.remove(0);
- if
(!SchemaBuilder.nestedUdts(userType).stream().allMatch(udts::containsKey))
+ if
(!AbstractSchemaBuilder.nestedUdts(userType).stream().allMatch(udts::containsKey))
{
// This UDT contains a nested user-defined type that has not
been parsed yet
// so re-add to the queue and parse later
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
index b565565b..306b8b8b 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
@@ -19,74 +19,21 @@
package org.apache.cassandra.spark.reader;
-import java.util.ArrayList;
import java.util.Collections;
-import java.util.HashMap;
-import java.util.HashSet;
-import java.util.Iterator;
-import java.util.List;
-import java.util.Map;
import java.util.Set;
import java.util.UUID;
-import java.util.function.Consumer;
import java.util.function.Function;
-import java.util.stream.Collectors;
import com.google.common.annotations.VisibleForTesting;
-import com.google.common.base.Preconditions;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.antlr.runtime.RecognitionException;
-import org.apache.cassandra.bridge.CassandraSchema;
-import org.apache.cassandra.bridge.CassandraTypesImplementation;
-import org.apache.cassandra.bridge.SchemaUpdater;
-import org.apache.cassandra.cql3.CQL3Type;
-import org.apache.cassandra.cql3.CQLFragmentParser;
-import org.apache.cassandra.cql3.CqlParser;
-import org.apache.cassandra.cql3.statements.schema.CreateTableStatement;
-import org.apache.cassandra.cql3.statements.schema.CreateTypeStatement;
-import org.apache.cassandra.db.Keyspace;
-import org.apache.cassandra.db.marshal.AbstractType;
-import org.apache.cassandra.db.marshal.CollectionType;
-import org.apache.cassandra.db.marshal.ListType;
-import org.apache.cassandra.db.marshal.MapType;
-import org.apache.cassandra.db.marshal.SetType;
-import org.apache.cassandra.db.marshal.TupleType;
-import org.apache.cassandra.db.marshal.UserType;
-import org.apache.cassandra.dht.IPartitioner;
-import org.apache.cassandra.schema.ColumnMetadata;
-import org.apache.cassandra.schema.KeyspaceMetadata;
-import org.apache.cassandra.schema.KeyspaceParams;
-import org.apache.cassandra.schema.Schema;
-import org.apache.cassandra.schema.TableId;
-import org.apache.cassandra.schema.TableMetadata;
-import org.apache.cassandra.schema.TableMetadataRef;
-import org.apache.cassandra.schema.Types;
import org.apache.cassandra.spark.data.CassandraTypes;
-import org.apache.cassandra.spark.data.CqlField;
import org.apache.cassandra.spark.data.CqlTable;
import org.apache.cassandra.spark.data.ReplicationFactor;
-import org.apache.cassandra.spark.data.complex.CqlFrozen;
-import org.apache.cassandra.spark.data.complex.CqlUdt;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
-import org.apache.cassandra.utils.Pair;
-import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
-public class SchemaBuilder
+public class SchemaBuilder extends AbstractSchemaBuilder
{
- private static final Logger LOGGER =
LoggerFactory.getLogger(SchemaBuilder.class);
-
- private final TableMetadata metadata;
- private final KeyspaceMetadata keyspaceMetadata;
- private final String createStmt;
- private final String keyspace;
- private final ReplicationFactor replicationFactor;
- private final CassandraTypes cassandraTypes;
- private final int indexCount;
- private final boolean enableCdc;
-
public SchemaBuilder(CqlTable table, Partitioner partitioner, boolean
enableCdc)
{
this(table, partitioner, null, enableCdc);
@@ -133,464 +80,7 @@ public class SchemaBuilder
int indexCount,
boolean enableCdc)
{
- this.createStmt = createStmt;
- this.keyspace = keyspace;
- this.replicationFactor = replicationFactor;
- this.cassandraTypes = new CassandraTypesImplementation();
- this.indexCount = indexCount;
- this.enableCdc = enableCdc;
-
- Pair<KeyspaceMetadata, TableMetadata> updated =
CassandraSchema.apply(schema ->
- updateSchema(schema,
- this.keyspace,
- udtStatementsProvider.apply(cassandraTypes),
- this.createStmt,
- partitioner,
- this.replicationFactor,
- tableId, enableCdc,
- this::validateColumnMetaData));
- this.keyspaceMetadata = updated.left;
- this.metadata = updated.right;
- }
-
- // Update schema with the given keyspace, table and udt.
- // It creates the corresponding metadata and opens instances for keyspace
and table, if needed.
- // At the end, it validates that the input keyspace and table both should
have metadata exist and instance opened.
- private static Pair<KeyspaceMetadata, TableMetadata> updateSchema(Schema
schema,
- String
keyspace,
-
Set<String> udtStatements,
- String
createStatement,
-
Partitioner partitioner,
-
ReplicationFactor replicationFactor,
- UUID
tableId,
- boolean
enableCdc,
-
Consumer<ColumnMetadata> columnValidator)
- {
- // Set up and open keyspace if needed
- IPartitioner cassPartitioner =
CassandraTypesImplementation.getPartitioner(partitioner);
- setupKeyspace(schema, keyspace, replicationFactor, cassPartitioner);
-
- // Set up and open table if needed, parse UDTs and include when
parsing table schema
- List<CreateTypeStatement.Raw> typeStatements = new
ArrayList<>(udtStatements.size());
- for (String udt : udtStatements)
- {
- try
- {
- typeStatements.add((CreateTypeStatement.Raw) CQLFragmentParser
- .parseAnyUnhandled(CqlParser::query, udt));
- }
- catch (RecognitionException exception)
- {
- LOGGER.error("Failed to parse type expression '{}'", udt);
- throw new IllegalStateException(exception);
- }
- }
- Types.RawBuilder typesBuilder = Types.rawBuilder(keyspace);
- for (CreateTypeStatement.Raw st : typeStatements)
- {
- st.addToRawBuilder(typesBuilder);
- }
- Types types = typesBuilder.build();
- CreateTableStatement.Raw createTable =
CQLFragmentParser.parseAny(CqlParser::createTableStatement,
-
createStatement,
-
"CREATE TABLE");
- // If the table already exists, the tableId should remain the same,
unless a non-null tableId is supplied
- TableMetadata maybeExistingTableMetadata =
schema.getTableMetadata(keyspace, createTable.table());
- if (maybeExistingTableMetadata != null && tableId == null)
- {
- tableId = maybeExistingTableMetadata.id.asUUID();
- }
-
- TableMetadata.Builder builder = createTable
- .keyspace(keyspace)
- .prepare(null)
- .builder(types)
- .partitioner(cassPartitioner);
-
- if (tableId != null)
- {
- builder.id(TableId.fromUUID(tableId));
- }
-
- TableMetadata tableMetadata = builder.build();
-
- if (tableMetadata.params.cdc != enableCdc)
- {
- tableMetadata = tableMetadata.unbuild()
- .params(tableMetadata.params.unbuild()
-
.cdc(enableCdc)
- .build())
- .build();
- }
-
- tableMetadata.columns().forEach(columnValidator);
- setupTableAndUdt(schema, keyspace, tableMetadata, types);
-
- return validateKeyspaceTable(schema, keyspace, tableMetadata.name);
- }
-
- private void validateColumnMetaData(@NotNull ColumnMetadata column)
- {
- validateType(column.type);
- }
-
- private void validateType(AbstractType<?> type)
- {
- validateType(type.asCQL3Type());
- }
-
- private void validateType(CQL3Type cqlType)
- {
- if (!(cqlType instanceof CQL3Type.Native)
- && !(cqlType instanceof CQL3Type.Collection)
- && !(cqlType instanceof CQL3Type.UserDefined)
- && !(cqlType instanceof CQL3Type.Tuple))
- {
- throw new UnsupportedOperationException("Only native, collection,
tuples or UDT data types are supported, "
- + "unsupported data type: "
+ cqlType.toString());
- }
-
- if (cqlType instanceof CQL3Type.Native)
- {
- CqlField.CqlType type =
cassandraTypes.parseType(cqlType.toString());
- if (!type.isSupported())
- {
- throw new UnsupportedOperationException(type.name() + " data
type is not supported");
- }
- }
- else if (cqlType instanceof CQL3Type.Collection)
- {
- // Validate collection inner types
- CQL3Type.Collection collection = (CQL3Type.Collection) cqlType;
- CollectionType<?> type = (CollectionType<?>) collection.getType();
- switch (type.kind)
- {
- case LIST:
- validateType(((ListType<?>) type).getElementsType());
- return;
- case SET:
- validateType(((SetType<?>) type).getElementsType());
- return;
- case MAP:
- validateType(((MapType<?, ?>) type).getKeysType());
- validateType(((MapType<?, ?>) type).getValuesType());
- return;
- default:
- // Do nothing
- }
- }
- else if (cqlType instanceof CQL3Type.Tuple)
- {
- CQL3Type.Tuple tuple = (CQL3Type.Tuple) cqlType;
- TupleType tupleType = (TupleType) tuple.getType();
- for (AbstractType<?> subType : tupleType.allTypes())
- {
- validateType(subType);
- }
- }
- else
- {
- // Validate UDT inner types
- UserType userType = (UserType) ((CQL3Type.UserDefined)
cqlType).getType();
- for (AbstractType<?> innerType : userType.fieldTypes())
- {
- validateType(innerType);
- }
- }
- }
-
- private static boolean keyspaceMetadataExists(Schema schema, String
keyspaceName)
- {
- return schema.getKeyspaceMetadata(keyspaceName) != null;
- }
-
- private static boolean tableMetadataExists(Schema schema, String
keyspaceName, String tableName)
- {
- KeyspaceMetadata ksMetadata = schema.getKeyspaceMetadata(keyspaceName);
- if (ksMetadata == null)
- {
- return false;
- }
-
- return ksMetadata.hasTable(tableName);
- }
-
- private static boolean keyspaceInstanceExists(Schema schema, String
keyspaceName)
- {
- return schema.getKeyspaceInstance(keyspaceName) != null;
- }
-
- private static boolean tableInstanceExists(Schema schema, String
keyspaceName, String tableName)
- {
- Keyspace keyspace = schema.getKeyspaceInstance(keyspaceName);
- if (keyspace == null)
- {
- return false;
- }
-
- try
- {
- keyspace.getColumnFamilyStore(tableName);
- }
- catch (IllegalArgumentException exception)
- {
- LOGGER.info("Table instance does not exist. keyspace={} table={}
existingCFS={}",
- keyspace, tableName, keyspace.getColumnFamilyStores());
- return false;
- }
- return true;
- }
-
- // Check whether keyspace metadata exists. Create keyspace metadata, if
not.
- // Check whether keyspace instance is opened. Open the keyspace, if not.
- // NOTE: It is possible that external code that just creates metadata, but
does not open the keyspace
- private static void setupKeyspace(Schema schema,
- String keyspaceName,
- ReplicationFactor replicationFactor,
- IPartitioner partitioner)
- {
- if (!keyspaceMetadataExists(schema, keyspaceName))
- {
- LOGGER.info("Setting up keyspace metadata in schema keyspace={}
rfStrategy={} partitioner={}",
- keyspaceName,
replicationFactor.getReplicationStrategy().name(), partitioner);
- KeyspaceMetadata keyspaceMetadata =
- KeyspaceMetadata.create(keyspaceName,
KeyspaceParams.create(true, rfToMap(replicationFactor)));
- SchemaUpdater.load(schema, keyspaceMetadata);
- }
-
- if (!keyspaceInstanceExists(schema, keyspaceName))
- {
- LOGGER.info("Setting up keyspace instance in schema keyspace={}
rfStrategy={} partitioner={}",
- keyspaceName,
replicationFactor.getReplicationStrategy().name(), partitioner);
- // Create keyspace instance and also initCf (cfs) for the table
- Keyspace.openWithoutSSTables(keyspaceName);
- }
- }
-
- // Check whether table metadata exists. Create table metadata, if not.
- // Check whether table instance is opened. Open/init the table, if not.
- // NOTE: It is possible that external code that just creates metadata, but
does not open the table
- private static void setupTableAndUdt(Schema schema,
- String keyspaceName,
- TableMetadata tableMetadata,
- Types userTypes)
- {
- String tableName = tableMetadata.name;
- KeyspaceMetadata keyspaceMetadata =
schema.getKeyspaceMetadata(keyspaceName);
- if (keyspaceMetadata == null)
- {
- LOGGER.error("Keyspace metadata does not exist. keyspace={}",
keyspaceName);
- throw new IllegalStateException("Keyspace metadata null for '" +
keyspaceName
- + "' when it should have been
initialized already");
- }
-
- if (!tableMetadataExists(schema, keyspaceName, tableName))
- {
- LOGGER.info("Setting up table metadata in schema keyspace={}
table={} partitioner={}",
- keyspaceName, tableName,
tableMetadata.partitioner.getClass().getName());
- keyspaceMetadata =
keyspaceMetadata.withSwapped(keyspaceMetadata.tables.with(tableMetadata));
- SchemaUpdater.load(schema, keyspaceMetadata, tableMetadata);
- }
-
- if (!tableMetadata.equals(schema.getTableMetadata(keyspaceName,
tableMetadata.name)))
- {
- // Schema of the table has changed so update it in the schema
- updateTableMetaData(schema, keyspaceName, tableMetadata);
- LOGGER.info("Table metadata changed schema keyspace={} table={}
partitioner={}",
- keyspaceName, tableName,
tableMetadata.partitioner.getClass().getName());
- }
-
- // The metadata of the table might not be the input tableMetadata.
Fetch the current to be safe.
- TableMetadata currentTable = schema.getTableMetadata(keyspaceName,
tableName);
- if (!tableInstanceExists(schema, keyspaceName, tableName))
- {
- LOGGER.info("Setting up table instance in schema keyspace={}
table={} partitioner={}",
- keyspaceName, tableName,
tableMetadata.partitioner.getClass().getName());
- if (keyspaceInstanceExists(schema, keyspaceName))
- {
- // initCf (cfs) in the opened keyspace
- schema.getKeyspaceInstance(keyspaceName)
- .initCf(TableMetadataRef.forOfflineTools(currentTable),
false);
- }
- else
- {
- // The keyspace has not yet opened, create/open keyspace
instance and also initCf (cfs) for the table
- Keyspace.openWithoutSSTables(keyspaceName);
- }
- }
-
- if (!userTypes.equals(Types.none()))
- {
- LOGGER.info("Setting up user types in schema keyspace={} types={}",
- keyspaceName, userTypes);
- // Update Schema instance with any user-defined types built
- keyspaceMetadata = keyspaceMetadata.withSwapped(userTypes);
- SchemaUpdater.load(schema, keyspaceMetadata, userTypes);
- }
- }
-
- private static void updateTableMetaData(Schema schema, String keyspace,
TableMetadata tableMetadata)
- {
- KeyspaceMetadata ks = schema.getKeyspaceMetadata(keyspace);
- ks = ks.withSwapped(ks.tables.withSwapped(tableMetadata));
- SchemaUpdater.load(schema, ks, tableMetadata);
- }
-
- private static Pair<KeyspaceMetadata, TableMetadata>
validateKeyspaceTable(Schema schema,
-
String keyspaceName,
-
String tableName)
- {
- Preconditions.checkState(keyspaceMetadataExists(schema, keyspaceName),
- "Keyspace metadata does not exist after
building schema. keyspace=%s",
- keyspaceName);
- Preconditions.checkState(keyspaceInstanceExists(schema, keyspaceName),
- "Keyspace instance is not opened after
building schema. keyspace=%s",
- keyspaceName);
- Preconditions.checkState(tableMetadataExists(schema, keyspaceName,
tableName),
- "Table metadata does not exist after building
schema. keyspace=%s table=%s",
- keyspaceName, tableName);
- Preconditions.checkState(tableInstanceExists(schema, keyspaceName,
tableName),
- "Table instance is not opened after building
schema. keyspace=%s table=%s",
- keyspaceName, tableName);
-
- // Validated above that keyspace and table, both exist and are opened
- KeyspaceMetadata keyspaceMetadata =
schema.getKeyspaceMetadata(keyspaceName);
- TableMetadata tableMetadata = schema.getTableMetadata(keyspaceName,
tableName);
- return Pair.create(keyspaceMetadata, tableMetadata);
- }
-
- public TableMetadata tableMetaData()
- {
- return metadata;
- }
-
- public String createStatement()
- {
- return createStmt;
- }
-
- public CqlTable build()
- {
- Map<String, CqlField.CqlUdt> udts = buildsUdts(keyspaceMetadata);
- List<CqlField> fields = buildFields(metadata,
udts).stream().sorted().collect(Collectors.toList());
- return new CqlTable(keyspace,
- metadata.name,
- createStmt,
- replicationFactor,
- fields,
- new HashSet<>(udts.values()),
- indexCount,
- enableCdc);
- }
-
- private Map<String, CqlField.CqlUdt> buildsUdts(KeyspaceMetadata
keyspaceMetadata)
- {
- List<UserType> userTypes = new ArrayList<>();
- keyspaceMetadata.types.forEach(userTypes::add);
- Map<String, CqlField.CqlUdt> udts = new HashMap<>(userTypes.size());
- while (!userTypes.isEmpty())
- {
- UserType userType = userTypes.remove(0);
- if
(!SchemaBuilder.nestedUdts(userType).stream().allMatch(udts::containsKey))
- {
- // This UDT contains a nested user-defined type that has not
been parsed yet
- // so re-add to the queue and parse later
- userTypes.add(userType);
- continue;
- }
- String name = userType.getNameAsString();
- CqlUdt.Builder builder = CqlUdt.builder(keyspaceMetadata.name,
name);
- for (int field = 0; field < userType.size(); field++)
- {
- builder.withField(userType.fieldName(field).toString(),
-
cassandraTypes.parseType(userType.fieldType(field).asCQL3Type().toString(),
udts));
- }
- udts.put(name, builder.build());
- }
-
- return udts;
- }
-
- /**
- * @param type an abstract type
- * @return a set of UDTs nested within the type parameter
- */
- private static Set<String> nestedUdts(AbstractType<?> type)
- {
- Set<String> result = new HashSet<>();
- nestedUdts(type, result, false);
- return result;
- }
-
- private static void nestedUdts(AbstractType<?> type, Set<String> udts,
boolean isNested)
- {
- if (type instanceof UserType)
- {
- if (isNested)
- {
- udts.add(((UserType) type).getNameAsString());
- }
- for (AbstractType<?> nestedType : ((UserType) type).fieldTypes())
- {
- nestedUdts(nestedType, udts, true);
- }
- }
- else if (type instanceof TupleType)
- {
- for (AbstractType<?> nestedType : ((TupleType) type).allTypes())
- {
- nestedUdts(nestedType, udts, true);
- }
- }
- else if (type instanceof SetType)
- {
- nestedUdts(((SetType<?>) type).getElementsType(), udts, true);
- }
- else if (type instanceof ListType)
- {
- nestedUdts(((ListType<?>) type).getElementsType(), udts, true);
- }
- else if (type instanceof MapType)
- {
- nestedUdts(((MapType<?, ?>) type).getKeysType(), udts, true);
- nestedUdts(((MapType<?, ?>) type).getValuesType(), udts, true);
- }
- }
-
- private List<CqlField> buildFields(TableMetadata metadata, Map<String,
CqlField.CqlUdt> udts)
- {
- Iterator<ColumnMetadata> it = metadata.allColumnsInSelectOrder();
- List<CqlField> result = new ArrayList<>();
- int position = 0;
- while (it.hasNext())
- {
- ColumnMetadata col = it.next();
- boolean isPartitionKey = col.isPartitionKey();
- boolean isClusteringColumn = col.isClusteringColumn();
- boolean isStatic = col.isStatic();
- String name = col.name.toString();
- CqlField.CqlType type = col.type.isUDT() ? udts.get(((UserType)
col.type).getNameAsString())
- :
cassandraTypes.parseType(col.type.asCQL3Type().toString(), udts);
- boolean isFrozen = col.type.isFreezable() &&
!col.type.isMultiCell();
- result.add(new CqlField(isPartitionKey,
- isClusteringColumn,
- isStatic,
- name,
- !(type instanceof CqlFrozen) && isFrozen ?
CqlFrozen.build(type) : type,
- position));
- position++;
- }
- return result;
- }
-
- static Map<String, String> rfToMap(ReplicationFactor replicationFactor)
- {
- Map<String, String> result = new
HashMap<>(replicationFactor.getOptions().size() + 1);
- result.put("class", "org.apache.cassandra.locator." +
replicationFactor.getReplicationStrategy().name());
- for (Map.Entry<String, Integer> entry :
replicationFactor.getOptions().entrySet())
- {
- result.put(entry.getKey(), Integer.toString(entry.getValue()));
- }
- return result;
+ super(createStmt, keyspace, replicationFactor, partitioner,
udtStatementsProvider,
+ tableId, indexCount, enableCdc);
}
}
diff --git a/gradle.properties b/gradle.properties
index aa3b798c..eaa11965 100644
--- a/gradle.properties
+++ b/gradle.properties
@@ -22,7 +22,7 @@ description=Apache Cassandra Analytics
analyticsJDKLevel=17
cassandra40Version=4.0.17
-cassandra50Version=5.0.5
+cassandra50Version=5.0.7
sidecarVersion=0.4.0
intellijVersion=9.0.4
junitVersion=5.10.2
diff --git a/gradlew b/gradlew
index 23d15a93..d9ea3c01 100755
--- a/gradlew
+++ b/gradlew
@@ -200,6 +200,12 @@ if "$cygwin" || "$msys" ; then
done
fi
+# We want to increase the file descriptor limit to the MaxFDLimit in MacOS,
which,
+# by default, is set to a lower limit than the actual system maximum. This
line is modified manually, if
+# producing a new gradle wrapper, remember to add the change back.
+if $darwin; then
+ GRADLE_OPTS="$GRADLE_OPTS \"-XX:-MaxFDLimit\"
\"-Dorg.gradle.jvmargs=-XX:-MaxFDLimit\""
+fi
# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to
pass JVM options to this script.
DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"'
diff --git a/scripts/build-dtest-jars.sh b/scripts/build-dtest-jars.sh
index fe7d4a96..6e2416be 100755
--- a/scripts/build-dtest-jars.sh
+++ b/scripts/build-dtest-jars.sh
@@ -40,12 +40,12 @@ else
#
# NOTE: The following branches need to stay in sync with the values in
build.gradle:
# ext.cassandraVersionEnumMap = ["4.0": "FOURZERO", "4.1": "FOURONE",
"5.0": "FIVEZERO"]
- # ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0":
"5.0.5"]
+ # ext.cassandraFullVersionMap = ["4.0": "4.0.17", "4.1": "4.1.4", "5.0":
"5.0.7"]
# NOTE: The following branches also need to remain in sync with
CassandraVersion.java
CANDIDATE_BRANCHES=(
"cassandra-4.0:cassandra-4.0.17"
"cassandra-4.1:99d9faeef57c9cf5240d11eac9db5b283e45a4f9"
- "cassandra-5.0:cassandra-5.0.5"
+ "cassandra-5.0:cassandra-5.0.7"
)
BRANCHES=( ${BRANCHES:-cassandra-4.0 cassandra-4.1 cassandra-5.0} )
echo ${BRANCHES[*]}
diff --git a/scripts/relocate-dtest-dependencies.pom
b/scripts/relocate-dtest-dependencies.pom
index 3108b6c4..bb0d5e2b 100644
--- a/scripts/relocate-dtest-dependencies.pom
+++ b/scripts/relocate-dtest-dependencies.pom
@@ -153,6 +153,7 @@
<artifact>*:*</artifact>
<excludes>
<exclude>**/Log4j2Plugins.dat</exclude>
+
<exclude>META-INF/versions/21/</exclude>
</excludes>
</filter>
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]