diff --git a/pom.xml b/pom.xml
index 0a52e73..b09eb9e 100644
--- a/pom.xml
+++ b/pom.xml
@@ -4,7 +4,7 @@
4.0.0
io.smartcat
cassandra-migration-tool
- 3.1.0.1-SNAPSHOT
+ 4.0.1-SNAPSHOT
jar
cassandra-migration-tool
@@ -56,8 +56,8 @@
UTF-8
1.8
1.8
- 3.1.0
- 3.0.0.1
+ 4.0.1
+ 4.3.1.0
18.0
1.7.7
4.12
@@ -76,8 +76,13 @@
- com.datastax.cassandra
- cassandra-driver-core
+ com.datastax.oss
+ java-driver-core
+ ${version.cassandra-driver}
+
+
+ com.datastax.oss
+ java-driver-query-builder
${version.cassandra-driver}
diff --git a/src/main/java/io/smartcat/migration/CassandraVersioner.java b/src/main/java/io/smartcat/migration/CassandraVersioner.java
index 9806f77..58adad0 100644
--- a/src/main/java/io/smartcat/migration/CassandraVersioner.java
+++ b/src/main/java/io/smartcat/migration/CassandraVersioner.java
@@ -3,12 +3,12 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import com.datastax.driver.core.ConsistencyLevel;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.Statement;
-import com.datastax.driver.core.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.DefaultConsistencyLevel;
+import com.datastax.oss.driver.api.core.cql.*;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder;
+import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal;
+import static com.datastax.oss.driver.api.querybuilder.relation.Relation.column;
/**
* Class responsible for version management.
@@ -32,13 +32,13 @@ public class CassandraVersioner {
+ String.format("PRIMARY KEY (%s, %s)", TYPE, VERSION)
+ String.format(") WITH CLUSTERING ORDER BY (%s DESC)", VERSION) + " AND COMMENT='Schema version';";
- private final Session session;
+ private final CqlSession session;
/**
* Create Cassandra versioner for active session.
* @param session Active Cassandra session
*/
- public CassandraVersioner(final Session session) {
+ public CassandraVersioner(final CqlSession session) {
this.session = session;
createSchemaVersion();
@@ -58,8 +58,9 @@ private void createSchemaVersion() {
* @return Database version for given type
*/
public int getCurrentVersion(final MigrationType type) {
- final Statement select = QueryBuilder.select().all().from(SCHEMA_VERSION_CF)
- .where(QueryBuilder.eq(TYPE, type.name())).limit(1).setConsistencyLevel(ConsistencyLevel.ALL);
+ final SimpleStatement select = QueryBuilder.selectFrom(SCHEMA_VERSION_CF).all()
+ .where(column(TYPE).isEqualTo(literal(type.name()))).limit(1)
+ .builder().setConsistencyLevel(DefaultConsistencyLevel.ALL).build();
final ResultSet result = session.execute(select);
final Row row = result.one();
@@ -73,10 +74,13 @@ public int getCurrentVersion(final MigrationType type) {
* @return Success of version update
*/
public boolean updateVersion(final Migration migration) {
- final Statement insert = QueryBuilder.insertInto(SCHEMA_VERSION_CF).value(TYPE, migration.getType().name())
- .value(VERSION, migration.getVersion()).value(TIMESTAMP, System.currentTimeMillis())
- .value(DESCRIPTION, migration.getDescription()).setConsistencyLevel(ConsistencyLevel.ALL);
-
+ final SimpleStatement insert = QueryBuilder.insertInto(SCHEMA_VERSION_CF)
+ .value(TYPE, literal(migration.getType().name()))
+ .value(VERSION, literal(migration.getVersion()))
+ .value(TIMESTAMP, literal(System.currentTimeMillis()))
+ .value(DESCRIPTION, literal(migration.getDescription()))
+ .builder().setConsistencyLevel(DefaultConsistencyLevel.ALL)
+ .build();
try {
session.execute(insert);
return true;
diff --git a/src/main/java/io/smartcat/migration/DataMigration.java b/src/main/java/io/smartcat/migration/DataMigration.java
index 124cd33..2e042fa 100644
--- a/src/main/java/io/smartcat/migration/DataMigration.java
+++ b/src/main/java/io/smartcat/migration/DataMigration.java
@@ -9,7 +9,7 @@ public abstract class DataMigration extends Migration {
* Creates new data migration.
* @param version Version of this data migration
*/
- public DataMigration(int version) {
+ protected DataMigration(int version) {
super(MigrationType.DATA, version);
}
}
diff --git a/src/main/java/io/smartcat/migration/Executor.java b/src/main/java/io/smartcat/migration/Executor.java
index d676e17..9a971b4 100644
--- a/src/main/java/io/smartcat/migration/Executor.java
+++ b/src/main/java/io/smartcat/migration/Executor.java
@@ -1,6 +1,6 @@
package io.smartcat.migration;
-import com.datastax.driver.core.Session;
+import com.datastax.oss.driver.api.core.CqlSession;
/**
* Executor is a class which executes all the migration for given session.
@@ -17,7 +17,7 @@ private Executor() {
* @param resources Migration resources collection
* @return Return success
*/
- public static boolean migrate(final Session session, final MigrationResources resources) {
+ public static boolean migrate(final CqlSession session, final MigrationResources resources) {
return MigrationEngine.withSession(session).migrate(resources);
}
diff --git a/src/main/java/io/smartcat/migration/Migration.java b/src/main/java/io/smartcat/migration/Migration.java
index 419a913..563da2a 100644
--- a/src/main/java/io/smartcat/migration/Migration.java
+++ b/src/main/java/io/smartcat/migration/Migration.java
@@ -1,8 +1,8 @@
package io.smartcat.migration;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.Statement;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.cql.ResultSet;
+import com.datastax.oss.driver.api.core.cql.Statement;
import io.smartcat.migration.exceptions.MigrationException;
import io.smartcat.migration.exceptions.SchemaAgreementException;
@@ -18,14 +18,14 @@ public abstract class Migration {
/**
* Active Cassandra session.
*/
- protected Session session;
+ protected CqlSession session;
/**
* Create new migration with provided type and version.
* @param type Migration type (SCHEMA or DATA)
* @param version Migration version
*/
- public Migration(final MigrationType type, final int version) {
+ protected Migration(final MigrationType type, final int version) {
this.type = type;
this.version = version;
}
@@ -34,7 +34,7 @@ public Migration(final MigrationType type, final int version) {
* Enables session injection into migration class.
* @param session Session object
*/
- public void setSession(final Session session) {
+ public void setSession(final CqlSession session) {
this.session = session;
}
@@ -78,7 +78,7 @@ protected void executeWithSchemaAgreement(Statement statement)
if (checkSchemaAgreement(result)) {
return;
}
- if (checkClusterSchemaAgreement()) {
+ if (checkSchemaAgreement()) {
return;
}
@@ -90,11 +90,11 @@ protected void executeWithSchemaAgreement(Statement statement)
* Whether the cluster had reached schema agreement after the execution of this query.
*
* After a successful schema-altering query (ex: creating a table), the driver will check if the cluster's nodes
- * agree on the new schema version. If not, it will keep retrying for a given delay (configurable via
- * {@link com.datastax.driver.core.Cluster.Builder#withMaxSchemaAgreementWaitSeconds(int)}).
+ * agree on the new schema version.
*
* If this method returns {@code false}, clients can call
- * {@link com.datastax.driver.core.Metadata#checkSchemaAgreement()} later to perform the check manually.
+ * {@link com.datastax.oss.driver.internal.core.cql.DefaultExecutionInfo#isSchemaInAgreement()}
+ * later to perform the check manually.
*
* Note that the schema agreement check is only performed for schema-altering queries For other query types, this
* method will always return {@code true}.
@@ -109,16 +109,12 @@ protected boolean checkSchemaAgreement(ResultSet resultSet) {
/**
* Checks whether hosts that are currently up agree on the schema definition.
*
- * This method performs a one-time check only, without any form of retry; therefore
- * {@link com.datastax.driver.core.Cluster.Builder#withMaxSchemaAgreementWaitSeconds(int)}
- * does not apply in this case.
- *
* @return {@code true} if all hosts agree on the schema; {@code false} if
* they don't agree, or if the check could not be performed
* (for example, if the control connection is down).
*/
- protected boolean checkClusterSchemaAgreement() {
- return this.session.getCluster().getMetadata().checkSchemaAgreement();
+ protected boolean checkSchemaAgreement() {
+ return this.session.checkSchemaAgreement();
}
@Override
diff --git a/src/main/java/io/smartcat/migration/MigrationEngine.java b/src/main/java/io/smartcat/migration/MigrationEngine.java
index 57c7340..e5edebc 100644
--- a/src/main/java/io/smartcat/migration/MigrationEngine.java
+++ b/src/main/java/io/smartcat/migration/MigrationEngine.java
@@ -4,7 +4,7 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import com.datastax.driver.core.Session;
+import com.datastax.oss.driver.api.core.CqlSession;
/**
* Migration engine wraps Migrator and provides DSL like API.
@@ -22,7 +22,7 @@ private MigrationEngine() {
* @param session Datastax driver session object
* @return migrator instance with versioner and session which can migrate resources
*/
- public static Migrator withSession(final Session session) {
+ public static Migrator withSession(final CqlSession session) {
return new Migrator(session);
}
@@ -30,14 +30,14 @@ public static Migrator withSession(final Session session) {
* Migrator handles migrations and errors.
*/
public static class Migrator {
- private final Session session;
+ private final CqlSession session;
private final CassandraVersioner versioner;
/**
* Create new Migrator with active Cassandra session.
* @param session Active Cassandra session
*/
- public Migrator(final Session session) {
+ public Migrator(final CqlSession session) {
this.session = session;
this.versioner = new CassandraVersioner(session);
}
diff --git a/src/main/java/io/smartcat/migration/SchemaMigration.java b/src/main/java/io/smartcat/migration/SchemaMigration.java
index f0df2a7..f27d319 100644
--- a/src/main/java/io/smartcat/migration/SchemaMigration.java
+++ b/src/main/java/io/smartcat/migration/SchemaMigration.java
@@ -9,7 +9,7 @@ public abstract class SchemaMigration extends Migration {
* Create new schema migration with provided version.
* @param version Version of this schema migration
*/
- public SchemaMigration(int version) {
+ protected SchemaMigration(int version) {
super(MigrationType.SCHEMA, version);
}
}
diff --git a/src/test/java/io/smartcat/migration/BaseTest.java b/src/test/java/io/smartcat/migration/BaseTest.java
index ade003a..d0c3695 100644
--- a/src/test/java/io/smartcat/migration/BaseTest.java
+++ b/src/test/java/io/smartcat/migration/BaseTest.java
@@ -1,28 +1,25 @@
package io.smartcat.migration;
-import com.datastax.driver.core.*;
+import com.datastax.oss.driver.api.core.CqlIdentifier;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.metadata.Metadata;
+import com.datastax.oss.driver.api.core.metadata.schema.KeyspaceMetadata;
-import java.util.ArrayList;
-import java.util.List;
+import java.util.Set;
public class BaseTest {
- public void truncateTables(final String keyspace, final Session session) {
- for (final String table : tables(keyspace, session)) {
+ public void truncateTables(final String keyspace, final CqlSession session) {
+ for (final CqlIdentifier table : tables(keyspace, session)) {
session.execute(String.format("TRUNCATE %s.%s;", keyspace, table));
}
}
- private List tables(final String keyspace, final Session session) {
- final List tables = new ArrayList<>();
- final Cluster cluster = session.getCluster();
- final Metadata meta = cluster.getMetadata();
- final KeyspaceMetadata keyspaceMeta = meta.getKeyspace(keyspace);
- for (final TableMetadata tableMeta : keyspaceMeta.getTables()) {
- tables.add(tableMeta.getName());
- }
+ private Set tables(final String keyspace, final CqlSession session) {
+ final Metadata meta = session.getMetadata();
+ final KeyspaceMetadata keyspaceMeta = meta.getKeyspace(keyspace).get();
- return tables;
+ return keyspaceMeta.getTables().keySet();
}
}
diff --git a/src/test/java/io/smartcat/migration/CassandraMetadataAnalyzer.java b/src/test/java/io/smartcat/migration/CassandraMetadataAnalyzer.java
index 9b35e36..2befac2 100644
--- a/src/test/java/io/smartcat/migration/CassandraMetadataAnalyzer.java
+++ b/src/test/java/io/smartcat/migration/CassandraMetadataAnalyzer.java
@@ -3,32 +3,34 @@
import static com.google.common.base.Preconditions.checkNotNull;
import static com.google.common.collect.Iterables.tryFind;
-import com.datastax.driver.core.ColumnMetadata;
-import com.datastax.driver.core.KeyspaceMetadata;
-import com.datastax.driver.core.Metadata;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.TableMetadata;
+import com.datastax.oss.driver.api.core.CqlIdentifier;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.metadata.Metadata;
+import com.datastax.oss.driver.api.core.metadata.schema.ColumnMetadata;
+import com.datastax.oss.driver.api.core.metadata.schema.KeyspaceMetadata;
+import com.datastax.oss.driver.api.core.metadata.schema.TableMetadata;
import com.google.common.base.Optional;
public class CassandraMetadataAnalyzer {
- private Session session;
+ private CqlSession session;
- public CassandraMetadataAnalyzer(Session session) {
+ public CassandraMetadataAnalyzer(CqlSession session) {
checkNotNull(session, "Session cannot be null");
- checkNotNull(session.getLoggedKeyspace(), "Session must be logged into a keyspace");
+ checkNotNull(session.getKeyspace(), "Session must be logged into a keyspace");
this.session = session;
}
public boolean columnExistInTable(String columnName, String tableName) {
TableMetadata table = getTableMetadata(this.session, tableName);
- Optional column = tryFind(table.getColumns(), new ColumnNameMatcher(columnName));
+ Optional column = tryFind(table.getColumns().values(), new ColumnNameMatcher(columnName));
return column.isPresent();
}
- private static TableMetadata getTableMetadata(Session session, String tableName) {
- Metadata metadata = session.getCluster().getMetadata();
- KeyspaceMetadata keyspaceMetadata = metadata.getKeyspace(session.getLoggedKeyspace());
- return keyspaceMetadata.getTable(tableName);
+ private static TableMetadata getTableMetadata(CqlSession session, String tableName) {
+ Metadata metadata = session.getMetadata();
+ CqlIdentifier keyspace = session.getKeyspace().get();
+ KeyspaceMetadata keyspaceMetadata = metadata.getKeyspace(keyspace).get();
+ return keyspaceMetadata.getTable(tableName).get();
}
}
diff --git a/src/test/java/io/smartcat/migration/CassandraVersionerTest.java b/src/test/java/io/smartcat/migration/CassandraVersionerTest.java
index ca8ee39..5662a81 100644
--- a/src/test/java/io/smartcat/migration/CassandraVersionerTest.java
+++ b/src/test/java/io/smartcat/migration/CassandraVersionerTest.java
@@ -6,32 +6,33 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.cql.ResultSet;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.core.cql.Statement;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mockito;
import org.mockito.stubbing.OngoingStubbing;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.Statement;
+
import io.smartcat.migration.migrations.schema.AddBookGenreFieldMigration;
public class CassandraVersionerTest {
private CassandraVersioner versioner;
- private Session session;
+ private CqlSession session;
private ResultSet versionResultSet;
@Before
- public void setUp() throws Exception {
- session = mock(Session.class);
+ public void setUp() {
+ session = mock(CqlSession.class);
versioner = new CassandraVersioner(session);
versionResultSet = mock(ResultSet.class);
}
@Test
- public void whenSchemaVersionTableIsEmptyThenCurrentVersionShouldBe0() throws Exception {
+ public void whenSchemaVersionTableIsEmptyThenCurrentVersionShouldBe0() {
expectRetrieveEmptyCurrentVersion();
int currentVersion = versioner.getCurrentVersion(SCHEMA);
@@ -40,7 +41,7 @@ public void whenSchemaVersionTableIsEmptyThenCurrentVersionShouldBe0() throws Ex
}
@Test
- public void whenSchemaVersionTableIsNotEmptyThenCurrentVersionShouldBeRetrievedFromTheTable() throws Exception {
+ public void whenSchemaVersionTableIsNotEmptyThenCurrentVersionShouldBeRetrievedFromTheTable() {
int expectedVersion = 1;
expectRetrieveCurrentVersion(expectedVersion);
@@ -51,7 +52,7 @@ public void whenSchemaVersionTableIsNotEmptyThenCurrentVersionShouldBeRetrievedF
}
@Test
- public void updateVersionSucess() throws Exception {
+ public void updateVersionSuccess() {
versioner.updateVersion(new AddBookGenreFieldMigration(1));
}
diff --git a/src/test/java/io/smartcat/migration/ColumnNameMatcher.java b/src/test/java/io/smartcat/migration/ColumnNameMatcher.java
index 6ca60e3..0ee9eba 100644
--- a/src/test/java/io/smartcat/migration/ColumnNameMatcher.java
+++ b/src/test/java/io/smartcat/migration/ColumnNameMatcher.java
@@ -1,6 +1,6 @@
package io.smartcat.migration;
-import com.datastax.driver.core.ColumnMetadata;
+import com.datastax.oss.driver.api.core.metadata.schema.ColumnMetadata;
import com.google.common.base.Predicate;
public class ColumnNameMatcher implements Predicate {
@@ -11,6 +11,6 @@ public ColumnNameMatcher(String columnName) {
}
public boolean apply(ColumnMetadata column) {
- return column.getName().equals(columnName);
+ return column.getName().asInternal().equals(columnName);
}
}
diff --git a/src/test/java/io/smartcat/migration/MigrationEngineBooksTest.java b/src/test/java/io/smartcat/migration/MigrationEngineBooksTest.java
index 66d04aa..4340135 100644
--- a/src/test/java/io/smartcat/migration/MigrationEngineBooksTest.java
+++ b/src/test/java/io/smartcat/migration/MigrationEngineBooksTest.java
@@ -1,6 +1,6 @@
package io.smartcat.migration;
-import static junit.framework.Assert.assertEquals;
+import static org.junit.Assert.assertEquals;
import org.cassandraunit.CQLDataLoader;
import org.cassandraunit.dataset.cql.ClassPathCQLDataSet;
@@ -10,13 +10,14 @@
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.net.InetSocketAddress;
-import com.datastax.driver.core.Cluster;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.Statement;
-import com.datastax.driver.core.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.CqlSessionBuilder;
+import com.datastax.oss.driver.api.core.cql.ResultSet;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import io.smartcat.migration.migrations.data.AddGenreMigration;
import io.smartcat.migration.migrations.data.InsertBooksMigration;
@@ -26,34 +27,31 @@ public class MigrationEngineBooksTest extends BaseTest {
private static final Logger LOGGER = LoggerFactory.getLogger(MigrationEngineBooksTest.class);
+ private static final String LOCAL_DATACENTER = "DC1";
private static final String CONTACT_POINT = "localhost";
private static final int PORT = 9142;
private static final String KEYSPACE = "migration_test_books";
private static final String CQL = "books.cql";
- private static Session session;
- private static Cluster cluster;
+ private static CqlSession session;
@BeforeClass
public static void init() throws Exception {
LOGGER.info("Starting embedded cassandra server");
EmbeddedCassandraServerHelper.startEmbeddedCassandra("another-cassandra.yaml");
- LOGGER.info("Connect to embedded db");
- cluster = Cluster.builder().addContactPoints(CONTACT_POINT).withPort(PORT).build();
- session = cluster.connect();
+ session = new CqlSessionBuilder()
+ .withLocalDatacenter(LOCAL_DATACENTER)
+ .addContactPoint(new InetSocketAddress(CONTACT_POINT, PORT))
+ .build();
LOGGER.info("Initialize keyspace");
- final CQLDataLoader cqlDataLoader = new CQLDataLoader(session);
- cqlDataLoader.load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
+ new CQLDataLoader(session).load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
}
@AfterClass
public static void tearDown() {
- if (cluster != null) {
- cluster.close();
- cluster = null;
- }
+ session.close();
}
@Test
@@ -76,7 +74,7 @@ public void test_data_migration() {
resources.addMigration(new AddGenreMigration(2));
MigrationEngine.withSession(session).migrate(resources);
- final Statement select = QueryBuilder.select().all().from("books");
+ final SimpleStatement select = QueryBuilder.selectFrom("books").all().build();
final ResultSet results = session.execute(select);
for (final Row row : results) {
diff --git a/src/test/java/io/smartcat/migration/MigrationEngineItemsTest.java b/src/test/java/io/smartcat/migration/MigrationEngineItemsTest.java
index a000893..dd5b6c9 100644
--- a/src/test/java/io/smartcat/migration/MigrationEngineItemsTest.java
+++ b/src/test/java/io/smartcat/migration/MigrationEngineItemsTest.java
@@ -1,9 +1,11 @@
package io.smartcat.migration;
-import static junit.framework.Assert.assertEquals;
+import static org.junit.Assert.assertEquals;
-import com.datastax.driver.core.*;
-import com.datastax.driver.core.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.CqlSessionBuilder;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder;
import io.smartcat.migration.migrations.data.InsertInitialItemsMigration;
import io.smartcat.migration.migrations.data.PopulateItemByNumberAndExternalIdMigration;
import io.smartcat.migration.migrations.schema.CreateItemByNumberAndExternalIdMigration;
@@ -17,19 +19,20 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.net.InetSocketAddress;
import java.util.List;
public class MigrationEngineItemsTest extends BaseTest {
private static final Logger LOGGER = LoggerFactory.getLogger(MigrationEngineItemsTest.class);
+ private static final String LOCAL_DATACENTER = "DC1";
private static final String CONTACT_POINT = "localhost";
private static final int PORT = 9142;
private static final String KEYSPACE = "migration_test_items";
private static final String CQL = "items.cql";
- private static Session session;
- private static Cluster cluster;
+ private static CqlSession session;
@BeforeClass
public static void init() throws Exception {
@@ -37,12 +40,13 @@ public static void init() throws Exception {
EmbeddedCassandraServerHelper.startEmbeddedCassandra("another-cassandra.yaml");
LOGGER.info("Connect to embedded db");
- cluster = Cluster.builder().addContactPoints(CONTACT_POINT).withPort(PORT).build();
- session = cluster.connect();
+ session = new CqlSessionBuilder()
+ .withLocalDatacenter(LOCAL_DATACENTER)
+ .addContactPoint(new InetSocketAddress(CONTACT_POINT, PORT))
+ .build();
LOGGER.info("Initialize keyspace");
- final CQLDataLoader cqlDataLoader = new CQLDataLoader(session);
- cqlDataLoader.load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
+ new CQLDataLoader(session).load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
}
@After
@@ -52,10 +56,7 @@ public void cleanUp() {
@AfterClass
public static void tearDown() {
- if (cluster != null) {
- cluster.close();
- cluster = null;
- }
+ EmbeddedCassandraServerHelper.cleanEmbeddedCassandra();
}
@Test
@@ -68,7 +69,8 @@ public void initial_insert_test() {
assertEquals(true, result);
- final List rows = session.execute(QueryBuilder.select().from("items_by_id")).all();
+ final List rows = session.execute(QueryBuilder.selectFrom("items_by_id").all().build())
+ .all();
assertEquals(count, rows.size());
}
@@ -84,7 +86,8 @@ public void test_migrations() {
assertEquals(true, result);
- final List rows = session.execute(QueryBuilder.select().from("items_by_number_external_id")).all();
+ final List rows = session.execute(QueryBuilder.selectFrom("items_by_number_external_id").all().build())
+ .all();
assertEquals(count, rows.size());
}
diff --git a/src/test/java/io/smartcat/migration/MigrationTest.java b/src/test/java/io/smartcat/migration/MigrationTest.java
index 2e66321..77cca95 100644
--- a/src/test/java/io/smartcat/migration/MigrationTest.java
+++ b/src/test/java/io/smartcat/migration/MigrationTest.java
@@ -1,7 +1,6 @@
package io.smartcat.migration;
-import static org.junit.Assert.assertFalse;
-import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.*;
import org.junit.Test;
@@ -12,7 +11,7 @@ public void test_equal_migrations() {
final Migration migration1 = new MigrationTestImplementation(MigrationType.SCHEMA, 1);
final Migration migration2 = new MigrationTestImplementation(MigrationType.SCHEMA, 1);
- assertTrue(migration1.equals(migration2));
+ assertEquals(migration1, migration2);
}
@Test
@@ -20,7 +19,7 @@ public void test_different_type_non_equal_migrations() {
final Migration migration1 = new MigrationTestImplementation(MigrationType.DATA, 1);
final Migration migration2 = new MigrationTestImplementation(MigrationType.SCHEMA, 1);
- assertFalse(migration1.equals(migration2));
+ assertNotEquals(migration1, migration2);
}
@Test
@@ -28,7 +27,7 @@ public void test_different_version_non_equal_migrations() {
final Migration migration1 = new MigrationTestImplementation(MigrationType.SCHEMA, 1);
final Migration migration2 = new MigrationTestImplementation(MigrationType.SCHEMA, 2);
- assertFalse(migration1.equals(migration2));
+ assertNotEquals(migration1, migration2);
}
public class MigrationTestImplementation extends Migration {
diff --git a/src/test/java/io/smartcat/migration/MigratorTest.java b/src/test/java/io/smartcat/migration/MigratorTest.java
index 2eed788..546e16d 100644
--- a/src/test/java/io/smartcat/migration/MigratorTest.java
+++ b/src/test/java/io/smartcat/migration/MigratorTest.java
@@ -2,9 +2,7 @@
import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.CoreMatchers.not;
-import static org.junit.Assert.assertFalse;
-import static org.junit.Assert.assertThat;
-import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.*;
import org.cassandraunit.CQLDataLoader;
import org.cassandraunit.dataset.cql.ClassPathCQLDataSet;
@@ -15,9 +13,10 @@
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.net.InetSocketAddress;
-import com.datastax.driver.core.Cluster;
-import com.datastax.driver.core.Session;
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.CqlSessionBuilder;
import io.smartcat.migration.MigrationEngine.Migrator;
import io.smartcat.migration.migrations.schema.AddBookGenreFieldMigration;
@@ -26,13 +25,13 @@
public class MigratorTest extends BaseTest {
private static final Logger LOGGER = LoggerFactory.getLogger(MigratorTest.class);
+ private static final String LOCAL_DATACENTER = "DC1";
private static final String CONTACT_POINT = "localhost";
private static final int PORT = 9142;
private static final String KEYSPACE = "migration_test_books";
private static final String CQL = "books.cql";
- private static Session session;
- private static Cluster cluster;
+ private static CqlSession session;
private CassandraVersioner versioner;
private Migrator migrator;
@@ -45,15 +44,16 @@ public static void init() throws Exception {
EmbeddedCassandraServerHelper.startEmbeddedCassandra("another-cassandra.yaml");
LOGGER.info("Connect to embedded db");
- cluster = Cluster.builder().addContactPoints(CONTACT_POINT).withPort(PORT).build();
- session = cluster.connect();
+ session = new CqlSessionBuilder()
+ .withLocalDatacenter(LOCAL_DATACENTER)
+ .addContactPoint(new InetSocketAddress(CONTACT_POINT, PORT))
+ .build();
}
@Before
- public void setUp() throws Exception {
+ public void setUp() {
LOGGER.info("Initialize keyspace");
- final CQLDataLoader cqlDataLoader = new CQLDataLoader(session);
- cqlDataLoader.load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
+ new CQLDataLoader(session).load(new ClassPathCQLDataSet(CQL, false, true, KEYSPACE));
versioner = new CassandraVersioner(session);
migrator = MigrationEngine.withSession(session);
@@ -62,14 +62,11 @@ public void setUp() throws Exception {
@AfterClass
public static void tearDown() {
- if (cluster != null) {
- cluster.close();
- cluster = null;
- }
+ session.close();
}
@Test
- public void executeOneMigration() throws Exception {
+ public void executeOneMigration() {
assertTableDoesntContainsColumns("books", "genre");
final MigrationResources resources = new MigrationResources();
@@ -81,7 +78,7 @@ public void executeOneMigration() throws Exception {
}
@Test
- public void executeTwoMigrations() throws Exception {
+ public void executeTwoMigrations() {
assertTableDoesntContainsColumns("books", "genre", "isbn");
final MigrationResources resources = new MigrationResources();
@@ -94,7 +91,7 @@ public void executeTwoMigrations() throws Exception {
}
@Test
- public void updateVersionAfterMigration() throws Exception {
+ public void updateVersionAfterMigration() {
int versionBeforeMigration = getCurrentVersion();
Migration migration = new AddBookGenreFieldMigration(1);
@@ -109,7 +106,7 @@ public void updateVersionAfterMigration() throws Exception {
}
@Test
- public void skipMigrationWithVersionOlderThanCurrentSchemaVersion() throws Exception {
+ public void skipMigrationWithVersionOlderThanCurrentSchemaVersion() {
Migration migrationWithNewerVersion = new AddBookGenreFieldMigration(2);
Migration migrationWithOlderVersion = new AddBookISBNFieldMigration(1);
@@ -125,7 +122,7 @@ public void skipMigrationWithVersionOlderThanCurrentSchemaVersion() throws Excep
}
@Test
- public void skipMigrationWithSameVersionThanCurrentSchemaVersion() throws Exception {
+ public void skipMigrationWithSameVersionThanCurrentSchemaVersion() {
int versionBeforeMigration = getCurrentVersion();
final MigrationResources resources = new MigrationResources();
diff --git a/src/test/java/io/smartcat/migration/migrations/data/AddGenreMigration.java b/src/test/java/io/smartcat/migration/migrations/data/AddGenreMigration.java
index d3024bc..81ffe52 100644
--- a/src/test/java/io/smartcat/migration/migrations/data/AddGenreMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/data/AddGenreMigration.java
@@ -1,11 +1,7 @@
package io.smartcat.migration.migrations.data;
-import com.datastax.driver.core.BoundStatement;
-import com.datastax.driver.core.PreparedStatement;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.Statement;
-import com.datastax.driver.core.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.cql.*;
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder;
import io.smartcat.migration.DataMigration;
import io.smartcat.migration.exceptions.MigrationException;
@@ -34,7 +30,7 @@ public void execute() throws MigrationException {
}
private void addGenreToBooks() {
- final Statement select = QueryBuilder.select().all().from("books");
+ final SimpleStatement select = QueryBuilder.selectFrom("books").all().build();
final ResultSet results = this.session.execute(select);
final PreparedStatement updateBookGenreStatement =
diff --git a/src/test/java/io/smartcat/migration/migrations/data/InsertBooksMigration.java b/src/test/java/io/smartcat/migration/migrations/data/InsertBooksMigration.java
index fba8b45..7292795 100644
--- a/src/test/java/io/smartcat/migration/migrations/data/InsertBooksMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/data/InsertBooksMigration.java
@@ -1,6 +1,6 @@
package io.smartcat.migration.migrations.data;
-import com.datastax.driver.core.PreparedStatement;
+import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import io.smartcat.migration.DataMigration;
import io.smartcat.migration.exceptions.MigrationException;
diff --git a/src/test/java/io/smartcat/migration/migrations/data/InsertInitialItemsMigration.java b/src/test/java/io/smartcat/migration/migrations/data/InsertInitialItemsMigration.java
index 289c396..80d7610 100644
--- a/src/test/java/io/smartcat/migration/migrations/data/InsertInitialItemsMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/data/InsertInitialItemsMigration.java
@@ -1,6 +1,6 @@
package io.smartcat.migration.migrations.data;
-import com.datastax.driver.core.PreparedStatement;
+import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import io.smartcat.migration.DataMigration;
import io.smartcat.migration.exceptions.MigrationException;
diff --git a/src/test/java/io/smartcat/migration/migrations/data/PopulateItemByNumberAndExternalIdMigration.java b/src/test/java/io/smartcat/migration/migrations/data/PopulateItemByNumberAndExternalIdMigration.java
index ccc9ba7..c82f40f 100644
--- a/src/test/java/io/smartcat/migration/migrations/data/PopulateItemByNumberAndExternalIdMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/data/PopulateItemByNumberAndExternalIdMigration.java
@@ -2,9 +2,10 @@
import java.util.List;
-import com.datastax.driver.core.PreparedStatement;
-import com.datastax.driver.core.Row;
-import com.datastax.driver.core.querybuilder.QueryBuilder;
+import com.datastax.oss.driver.api.core.cql.PreparedStatement;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder;
import io.smartcat.migration.DataMigration;
import io.smartcat.migration.exceptions.MigrationException;
@@ -27,10 +28,11 @@ public void execute() throws MigrationException {
session.prepare(
"INSERT INTO items_by_number_external_id (id, number, external_id) VALUES (?, ?, ?);");
- final List rows = session.execute(QueryBuilder.select().from("items_by_id").setFetchSize(1000)).all();
+ final SimpleStatement select = QueryBuilder.selectFrom("items_by_id").all().limit(1000).build();
+ final List rows = session.execute(select).all();
for (Row row : rows) {
session.execute(
- preparedStatement.bind(row.getUUID("id"), row.getString("number"), row.getUUID("external_id")));
+ preparedStatement.bind(row.getUuid("id"), row.getString("number"), row.getUuid("external_id")));
}
} catch (final Exception e) {
throw new MigrationException("Failed to execute PopulateItemByNumberAndExternalId migration", e);
diff --git a/src/test/java/io/smartcat/migration/migrations/schema/AddBookGenreFieldMigration.java b/src/test/java/io/smartcat/migration/migrations/schema/AddBookGenreFieldMigration.java
index 25c7df9..bb7f70c 100644
--- a/src/test/java/io/smartcat/migration/migrations/schema/AddBookGenreFieldMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/schema/AddBookGenreFieldMigration.java
@@ -1,6 +1,6 @@
package io.smartcat.migration.migrations.schema;
-import com.datastax.driver.core.SimpleStatement;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import io.smartcat.migration.exceptions.MigrationException;
import io.smartcat.migration.SchemaMigration;
@@ -24,7 +24,7 @@ public void execute() throws MigrationException {
try {
final String alterBooksAddGenreCQL = "ALTER TABLE books ADD genre text;";
- executeWithSchemaAgreement(new SimpleStatement(alterBooksAddGenreCQL));
+ executeWithSchemaAgreement(SimpleStatement.newInstance(alterBooksAddGenreCQL));
} catch (final Exception e) {
throw new MigrationException("Failed to execute AddBookGenreField migration", e);
diff --git a/src/test/java/io/smartcat/migration/migrations/schema/AddBookISBNFieldMigration.java b/src/test/java/io/smartcat/migration/migrations/schema/AddBookISBNFieldMigration.java
index 5be8243..db3ec7b 100644
--- a/src/test/java/io/smartcat/migration/migrations/schema/AddBookISBNFieldMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/schema/AddBookISBNFieldMigration.java
@@ -1,6 +1,6 @@
package io.smartcat.migration.migrations.schema;
-import com.datastax.driver.core.SimpleStatement;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import io.smartcat.migration.SchemaMigration;
import io.smartcat.migration.exceptions.MigrationException;
@@ -24,7 +24,7 @@ public void execute() throws MigrationException {
try {
final String alterBooksAddISBNCQL = "ALTER TABLE books ADD isbn text;";
- executeWithSchemaAgreement(new SimpleStatement(alterBooksAddISBNCQL));
+ executeWithSchemaAgreement(SimpleStatement.newInstance(alterBooksAddISBNCQL));
} catch (final Exception e) {
throw new MigrationException("Failed to execute AddBookISBNField migration", e);
diff --git a/src/test/java/io/smartcat/migration/migrations/schema/CreateItemByNumberAndExternalIdMigration.java b/src/test/java/io/smartcat/migration/migrations/schema/CreateItemByNumberAndExternalIdMigration.java
index 19d220a..b7193d7 100644
--- a/src/test/java/io/smartcat/migration/migrations/schema/CreateItemByNumberAndExternalIdMigration.java
+++ b/src/test/java/io/smartcat/migration/migrations/schema/CreateItemByNumberAndExternalIdMigration.java
@@ -1,6 +1,6 @@
package io.smartcat.migration.migrations.schema;
-import com.datastax.driver.core.SimpleStatement;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
import io.smartcat.migration.SchemaMigration;
import io.smartcat.migration.exceptions.MigrationException;
@@ -25,7 +25,7 @@ public void execute() throws MigrationException {
"external_id uuid," +
"PRIMARY KEY ((number, external_id))" +
") WITH COMMENT='Items by item number and external id';";
- executeWithSchemaAgreement(new SimpleStatement(statement));
+ executeWithSchemaAgreement(SimpleStatement.newInstance(statement));
} catch (final Exception e) {
throw new MigrationException("Failed to execute CreateItemsByNumberAndExternalIdMigration migration", e);
diff --git a/src/test/resources/another-cassandra.yaml b/src/test/resources/another-cassandra.yaml
index f336e03..c6d2ce6 100644
--- a/src/test/resources/another-cassandra.yaml
+++ b/src/test/resources/another-cassandra.yaml
@@ -74,6 +74,8 @@ data_file_directories:
# commit log
commitlog_directory: target/embeddedCassandra/commitlog
+cdc_raw_directory: target/embeddedCassandra/cdc
+
# Maximum size of the key cache in memory.
#
# Each key cache hit saves 1 seek and each row cache hit saves 2 seeks at the