From bfc31d82f60b987a572e46545a262af1c94bf763 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Boles=C5=82aw=20Hangiel?= Date: Wed, 16 Dec 2020 12:45:55 +0100 Subject: [PATCH] 4.X cassandra driver update --- pom.xml | 15 ++++--- .../migration/CassandraVersioner.java | 32 ++++++++------- .../io/smartcat/migration/DataMigration.java | 2 +- .../java/io/smartcat/migration/Executor.java | 4 +- .../java/io/smartcat/migration/Migration.java | 28 ++++++------- .../smartcat/migration/MigrationEngine.java | 8 ++-- .../smartcat/migration/SchemaMigration.java | 2 +- .../java/io/smartcat/migration/BaseTest.java | 25 ++++++------ .../migration/CassandraMetadataAnalyzer.java | 28 ++++++------- .../migration/CassandraVersionerTest.java | 21 +++++----- .../smartcat/migration/ColumnNameMatcher.java | 4 +- .../migration/MigrationEngineBooksTest.java | 36 ++++++++--------- .../migration/MigrationEngineItemsTest.java | 33 +++++++++------- .../io/smartcat/migration/MigrationTest.java | 9 ++--- .../io/smartcat/migration/MigratorTest.java | 39 +++++++++---------- .../migrations/data/AddGenreMigration.java | 10 ++--- .../migrations/data/InsertBooksMigration.java | 2 +- .../data/InsertInitialItemsMigration.java | 2 +- ...ateItemByNumberAndExternalIdMigration.java | 12 +++--- .../schema/AddBookGenreFieldMigration.java | 4 +- .../schema/AddBookISBNFieldMigration.java | 4 +- ...ateItemByNumberAndExternalIdMigration.java | 4 +- src/test/resources/another-cassandra.yaml | 2 + 23 files changed, 164 insertions(+), 162 deletions(-) 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