From 4d73835bb49316993ae5c784a17173d675af873e Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Wed, 4 Mar 2020 10:08:47 +0200 Subject: [PATCH] Cassansra to PostgreSQL migration. --- .../install/ThingsboardInstallService.java | 4 +- .../CassandraEntitiesToSqlMigrateService.java | 58 ++++- .../install/migrate/CassandraToSqlColumn.java | 70 +++--- .../migrate/CassandraToSqlColumnData.java | 64 +++++ .../migrate/CassandraToSqlColumnType.java | 1 + .../install/migrate/CassandraToSqlTable.java | 238 ++++++++++++++++-- 6 files changed, 366 insertions(+), 69 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnData.java diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index e0daf43bf7..b024d6a996 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -80,8 +80,8 @@ public class ThingsboardInstallService { if ("2.5.0-cassandra".equals(upgradeFromVersion)) { log.info("Migrating ThingsBoard entities data from cassandra to SQL database ..."); entitiesMigrateService.migrate(); - log.info("Updating system data..."); - systemDataLoaderService.updateSystemWidgets(); + // log.info("Updating system data..."); + // systemDataLoaderService.updateSystemWidgets(); } else { switch (upgradeFromVersion) { case "1.2.3": //NOSONAR, Need to execute gradual upgrade starting from upgradeFromVersion diff --git a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraEntitiesToSqlMigrateService.java b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraEntitiesToSqlMigrateService.java index eb3373fe7f..8719df5f36 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraEntitiesToSqlMigrateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraEntitiesToSqlMigrateService.java @@ -21,6 +21,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Profile; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.dao.cassandra.CassandraCluster; import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.service.install.EntityDatabaseSchemaService; @@ -35,6 +36,7 @@ import static org.thingsboard.server.service.install.migrate.CassandraToSqlColum import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.doubleColumn; import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.enumToIntColumn; import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.idColumn; +import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.jsonColumn; import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.stringColumn; @Service @@ -64,7 +66,7 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ log.info("Performing migration of entities data from cassandra to SQL database ..."); entityDatabaseSchemaService.createDatabaseSchema(); try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { - conn.setAutoCommit(true); + conn.setAutoCommit(false); for (CassandraToSqlTable table: tables) { table.migrateToSql(cluster.getSession(), conn); } @@ -75,7 +77,7 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ } private static List tables = Arrays.asList( - new CassandraToSqlTable("admin_settings", + new CassandraToSqlTable("admin_settings", idColumn("id"), stringColumn("key"), stringColumn("json_value")), @@ -102,7 +104,17 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ stringColumn("type"), stringColumn("label"), stringColumn("search_text"), - stringColumn("additional_info")), + stringColumn("additional_info")) { + @Override + protected boolean onConstraintViolation(List batchData, + CassandraToSqlColumnData[] data, String constraint) { + if (constraint.equalsIgnoreCase("asset_name_unq_key")) { + this.handleUniqueNameViolation(data, "asset"); + return true; + } + return super.onConstraintViolation(batchData, data, constraint); + } + }, new CassandraToSqlTable("audit_log_by_tenant_id", "audit_log", idColumn("id"), idColumn("tenant_id"), @@ -125,7 +137,7 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ stringColumn("str_v"), bigintColumn("long_v"), doubleColumn("dbl_v"), - stringColumn("json_v"), + jsonColumn("json_v"), bigintColumn("last_update_ts")), new CassandraToSqlTable("component_descriptor", idColumn("id"), @@ -165,7 +177,17 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ stringColumn("type"), stringColumn("label"), stringColumn("search_text"), - stringColumn("additional_info")), + stringColumn("additional_info")) { + @Override + protected boolean onConstraintViolation(List batchData, + CassandraToSqlColumnData[] data, String constraint) { + if (constraint.equalsIgnoreCase("device_name_unq_key")) { + this.handleUniqueNameViolation(data, "device"); + return true; + } + return super.onConstraintViolation(batchData, data, constraint); + } + }, new CassandraToSqlTable("device_credentials", idColumn("id"), idColumn("device_id"), @@ -197,7 +219,17 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ stringColumn("authority"), stringColumn("first_name"), stringColumn("last_name"), - stringColumn("additional_info")), + stringColumn("additional_info")) { + @Override + protected boolean onConstraintViolation(List batchData, + CassandraToSqlColumnData[] data, String constraint) { + if (constraint.equalsIgnoreCase("tb_user_email_key")) { + this.handleUniqueEmailViolation(data); + return true; + } + return super.onConstraintViolation(batchData, data, constraint); + } + }, new CassandraToSqlTable("tenant", idColumn("id"), stringColumn("title"), @@ -218,7 +250,19 @@ public class CassandraEntitiesToSqlMigrateService implements EntitiesMigrateServ booleanColumn("enabled"), stringColumn("password"), stringColumn("activate_token"), - stringColumn("reset_token")), + stringColumn("reset_token")) { + @Override + protected boolean onConstraintViolation(List batchData, + CassandraToSqlColumnData[] data, String constraint) { + if (constraint.equalsIgnoreCase("user_credentials_user_id_key")) { + String id = UUIDConverter.fromString(this.getColumnData(data, "id").getValue()).toString(); + log.warn("Found user credentials record with duplicate user_id [id:[{}]]. Record will be ignored!", id); + this.ignoreRecord(batchData, data); + return true; + } + return super.onConstraintViolation(batchData, data, constraint); + } + }, new CassandraToSqlTable("widget_type", idColumn("id"), idColumn("tenant_id"), diff --git a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java index eaddc320fc..6ade15a403 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java +++ b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java @@ -22,14 +22,21 @@ import org.thingsboard.server.common.data.UUIDConverter; import java.sql.PreparedStatement; import java.sql.SQLException; import java.sql.Types; +import java.util.regex.Pattern; @Data public class CassandraToSqlColumn { + private static final ThreadLocal PATTERN_THREAD_LOCAL = ThreadLocal.withInitial(() -> Pattern.compile(String.valueOf(Character.MIN_VALUE))); + private static final String EMPTY_STR = ""; + + private int index; + private int sqlIndex; private String cassandraColumnName; private String sqlColumnName; private CassandraToSqlColumnType type; private int sqlType; + private int size; private Class enumClass; public static CassandraToSqlColumn idColumn(String name) { @@ -56,6 +63,10 @@ public class CassandraToSqlColumn { return new CassandraToSqlColumn(name, CassandraToSqlColumnType.BOOLEAN); } + public static CassandraToSqlColumn jsonColumn(String name) { + return new CassandraToSqlColumn(name, CassandraToSqlColumnType.JSON); + } + public static CassandraToSqlColumn enumToIntColumn(String name, Class enumClass) { return new CassandraToSqlColumn(name, CassandraToSqlColumnType.ENUM_TO_INT, enumClass); } @@ -82,36 +93,9 @@ public class CassandraToSqlColumn { this.sqlColumnName = sqlColumnName; this.type = type; this.enumClass = enumClass; - switch (this.type) { - case ID: - case STRING: - this.sqlType = Types.VARCHAR; - break; - case DOUBLE: - this.sqlType = Types.DOUBLE; - break; - case INTEGER: - case ENUM_TO_INT: - this.sqlType = Types.INTEGER; - break; - case FLOAT: - this.sqlType = Types.FLOAT; - break; - case BIGINT: - this.sqlType = Types.BIGINT; - break; - case BOOLEAN: - this.sqlType = Types.BOOLEAN; - break; - } } - public void prepareColumnValue(Row row, PreparedStatement sqlInsertStatement, int index) throws SQLException { - String value = this.getColumnValue(row, index); - this.setColumnValue(sqlInsertStatement, index, value); - } - - private String getColumnValue(Row row, int index) { + public String getColumnValue(Row row) { if (row.isNull(index)) { return null; } else { @@ -129,46 +113,56 @@ public class CassandraToSqlColumn { case BOOLEAN: return Boolean.toString(row.getBool(index)); case STRING: + case JSON: case ENUM_TO_INT: default: - return row.getString(index); + String value = row.getString(index); + return this.replaceNullChars(value); } } } - private void setColumnValue(PreparedStatement sqlInsertStatement, int index, String value) throws SQLException { + public void setColumnValue(PreparedStatement sqlInsertStatement, String value) throws SQLException { if (value == null) { - sqlInsertStatement.setNull(index, this.sqlType); + sqlInsertStatement.setNull(this.sqlIndex, this.sqlType); } else { switch (this.type) { case DOUBLE: - sqlInsertStatement.setDouble(index, Double.parseDouble(value)); + sqlInsertStatement.setDouble(this.sqlIndex, Double.parseDouble(value)); break; case INTEGER: - sqlInsertStatement.setInt(index, Integer.parseInt(value)); + sqlInsertStatement.setInt(this.sqlIndex, Integer.parseInt(value)); break; case FLOAT: - sqlInsertStatement.setFloat(index, Float.parseFloat(value)); + sqlInsertStatement.setFloat(this.sqlIndex, Float.parseFloat(value)); break; case BIGINT: - sqlInsertStatement.setLong(index, Long.parseLong(value)); + sqlInsertStatement.setLong(this.sqlIndex, Long.parseLong(value)); break; case BOOLEAN: - sqlInsertStatement.setBoolean(index, Boolean.parseBoolean(value)); + sqlInsertStatement.setBoolean(this.sqlIndex, Boolean.parseBoolean(value)); break; case ENUM_TO_INT: Enum enumVal = Enum.valueOf(this.enumClass, value); int intValue = enumVal.ordinal(); - sqlInsertStatement.setInt(index, intValue); + sqlInsertStatement.setInt(this.sqlIndex, intValue); break; + case JSON: case STRING: case ID: default: - sqlInsertStatement.setString(index, value); + sqlInsertStatement.setString(this.sqlIndex, value); break; } } } + private String replaceNullChars(String strValue) { + if (strValue != null) { + return PATTERN_THREAD_LOCAL.get().matcher(strValue).replaceAll(EMPTY_STR); + } + return strValue; + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnData.java b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnData.java new file mode 100644 index 0000000000..6f7052e56f --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnData.java @@ -0,0 +1,64 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.service.install.migrate; + +import lombok.Data; + +@Data +public class CassandraToSqlColumnData { + + private String value; + private String originalValue; + private int constraintCounter = 0; + + public CassandraToSqlColumnData(String value) { + this.value = value; + this.originalValue = value; + } + + public int nextContraintCounter() { + return ++constraintCounter; + } + + public String getNextConstraintStringValue(CassandraToSqlColumn column) { + int counter = this.nextContraintCounter(); + String newValue = this.originalValue + counter; + int overflow = newValue.length() - column.getSize(); + if (overflow > 0) { + newValue = this.originalValue.substring(0, this.originalValue.length()-overflow) + counter; + } + return newValue; + } + + public String getNextConstraintEmailValue(CassandraToSqlColumn column) { + int counter = this.nextContraintCounter(); + String[] emailValues = this.originalValue.split("@"); + String newValue = emailValues[0] + "+" + counter + "@" + emailValues[1]; + int overflow = newValue.length() - column.getSize(); + if (overflow > 0) { + newValue = emailValues[0].substring(0, emailValues[0].length()-overflow) + "+" + counter + "@" + emailValues[1]; + } + return newValue; + } + + public String getLogValue() { + if (this.value != null && this.value.length() > 255) { + return this.value.substring(0, 255) + "...[truncated " + (this.value.length() - 255) + " symbols]"; + } + return this.value; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnType.java b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnType.java index 25972d1bb3..f97b01ae77 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnType.java +++ b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumnType.java @@ -23,5 +23,6 @@ public enum CassandraToSqlColumnType { BIGINT, BOOLEAN, STRING, + JSON, ENUM_TO_INT } diff --git a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlTable.java b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlTable.java index 866ff7f5dc..2f806f0832 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlTable.java +++ b/application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlTable.java @@ -22,66 +22,255 @@ import com.datastax.driver.core.SimpleStatement; import com.datastax.driver.core.Statement; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.hibernate.exception.ConstraintViolationException; +import org.hibernate.internal.util.JdbcExceptionHelper; +import org.postgresql.util.PSQLException; +import org.thingsboard.server.common.data.UUIDConverter; +import org.thingsboard.server.dao.exception.DataValidationException; +import java.sql.BatchUpdateException; import java.sql.Connection; +import java.sql.DatabaseMetaData; import java.sql.PreparedStatement; import java.sql.SQLException; +import java.util.ArrayList; import java.util.Arrays; import java.util.Iterator; import java.util.List; +import java.util.Optional; +import java.util.stream.Collectors; @Data @Slf4j public class CassandraToSqlTable { + private static final int DEFAULT_BATCH_SIZE = 10000; + private String cassandraCf; private String sqlTableName; private List columns; + private int batchSize = DEFAULT_BATCH_SIZE; + + private PreparedStatement sqlInsertStatement; + public CassandraToSqlTable(String tableName, CassandraToSqlColumn... columns) { - this(tableName, tableName, columns); + this(tableName, tableName, DEFAULT_BATCH_SIZE, columns); + } + + public CassandraToSqlTable(String tableName, String sqlTableName, CassandraToSqlColumn... columns) { + this(tableName, sqlTableName, DEFAULT_BATCH_SIZE, columns); + } + + public CassandraToSqlTable(String tableName, int batchSize, CassandraToSqlColumn... columns) { + this(tableName, tableName, batchSize, columns); } - public CassandraToSqlTable(String cassandraCf, String sqlTableName, CassandraToSqlColumn... columns) { + public CassandraToSqlTable(String cassandraCf, String sqlTableName, int batchSize, CassandraToSqlColumn... columns) { this.cassandraCf = cassandraCf; this.sqlTableName = sqlTableName; + this.batchSize = batchSize; this.columns = Arrays.asList(columns); + for (int i=0;i iter = rs.iterator(); int rowCounter = 0; - while (iter.hasNext()) { + List batchData; + boolean hasNext; + do { + batchData = this.extractBatchData(iter); + hasNext = batchData.size() == this.batchSize; + this.batchInsert(batchData, conn); + rowCounter += batchData.size(); + log.info("[{}] {} records migrated so far...", this.sqlTableName, rowCounter); + } while (hasNext); + this.sqlInsertStatement.close(); + log.info("[{}] {} total records migrated.", this.sqlTableName, rowCounter); + log.info("[{}] Finished migration data from cassandra '{}' Column Family to '{}' SQL table.", + this.sqlTableName, this.cassandraCf, this.sqlTableName); + } + + private List extractBatchData(Iterator iter) { + List batchData = new ArrayList<>(); + while (iter.hasNext() && batchData.size() < this.batchSize) { Row row = iter.next(); if (row != null) { - this.migrateRowToSql(row, sqlInsertStatement); - rowCounter++; - if (rowCounter % 100 == 0) { - sqlInsertStatement.executeBatch(); - log.info("{} records migrated so far...", rowCounter); - } + CassandraToSqlColumnData[] data = this.extractRowData(row); + batchData.add(data); } } - if (rowCounter % 100 > 0) { - sqlInsertStatement.executeBatch(); + return batchData; + } + + private CassandraToSqlColumnData[] extractRowData(Row row) { + CassandraToSqlColumnData[] data = new CassandraToSqlColumnData[this.columns.size()]; + for (CassandraToSqlColumn column: this.columns) { + String value = column.getColumnValue(row); + data[column.getIndex()] = new CassandraToSqlColumnData(value); } - sqlInsertStatement.close(); - log.info("{} total records migrated.", rowCounter); - log.info("Finished migration data from cassandra '{}' Column Family to '{}' SQL table.", this.cassandraCf, this.sqlTableName); + return this.validateColumnData(data); } - private void migrateRowToSql(Row row, PreparedStatement sqlInsertStatement) throws SQLException { - for (int i=0; i column.getSize()) { + log.warn("[{}] Value size [{}] exceeds maximum size [{}] of column [{}] and will be truncated!", + this.sqlTableName, + value.length(), column.getSize(), column.getSqlColumnName()); + log.warn("[{}] Affected data:\n{}", this.sqlTableName, this.dataToString(data)); + value = value.substring(0, column.getSize()); + columnData.setOriginalValue(value); + columnData.setValue(value); + } + } + } + return data; + } + + private void batchInsert(List batchData, Connection conn) throws SQLException { + boolean retry = false; + for (CassandraToSqlColumnData[] data : batchData) { + for (CassandraToSqlColumn column: this.columns) { + column.setColumnValue(this.sqlInsertStatement, data[column.getIndex()].getValue()); + } + try { + this.sqlInsertStatement.executeUpdate(); + } catch (SQLException e) { + if (this.handleInsertException(batchData, data, conn, e)) { + retry = true; + break; + } else { + throw e; + } + } + } + if (retry) { + this.batchInsert(batchData, conn); + } else { + conn.commit(); + } + } + + private boolean handleInsertException(List batchData, + CassandraToSqlColumnData[] data, + Connection conn, SQLException ex) throws SQLException { + conn.commit(); + String constraint = extractConstraintName(ex).orElse(null); + if (constraint != null) { + if (this.onConstraintViolation(batchData, data, constraint)) { + return true; + } else { + log.error("[{}] Unhandled constraint violation [{}] during insert!", this.sqlTableName, constraint); + log.error("[{}] Affected data:\n{}", this.sqlTableName, this.dataToString(data)); + } + } else { + log.error("[{}] Unhandled exception during insert!", this.sqlTableName); + log.error("[{}] Affected data:\n{}", this.sqlTableName, this.dataToString(data)); } - sqlInsertStatement.addBatch(); + return false; + } + + private String dataToString(CassandraToSqlColumnData[] data) { + StringBuffer stringData = new StringBuffer("{\n"); + for (int i=0;i batchData, + CassandraToSqlColumnData[] data, String constraint) { + return false; + } + + protected void handleUniqueNameViolation(CassandraToSqlColumnData[] data, String entityType) { + CassandraToSqlColumn nameColumn = this.getColumn("name"); + CassandraToSqlColumn searchTextColumn = this.getColumn("search_text"); + CassandraToSqlColumnData nameColumnData = data[nameColumn.getIndex()]; + CassandraToSqlColumnData searchTextColumnData = data[searchTextColumn.getIndex()]; + String prevName = nameColumnData.getValue(); + String newName = nameColumnData.getNextConstraintStringValue(nameColumn); + nameColumnData.setValue(newName); + searchTextColumnData.setValue(searchTextColumnData.getNextConstraintStringValue(searchTextColumn)); + String id = UUIDConverter.fromString(this.getColumnData(data, "id").getValue()).toString(); + log.warn("Found {} with duplicate name [id:[{}]]. Attempting to rename {} from '{}' to '{}'...", entityType, id, entityType, prevName, newName); + } + + protected void handleUniqueEmailViolation(CassandraToSqlColumnData[] data) { + CassandraToSqlColumn emailColumn = this.getColumn("email"); + CassandraToSqlColumn searchTextColumn = this.getColumn("search_text"); + CassandraToSqlColumnData emailColumnData = data[emailColumn.getIndex()]; + CassandraToSqlColumnData searchTextColumnData = data[searchTextColumn.getIndex()]; + String prevEmail = emailColumnData.getValue(); + String newEmail = emailColumnData.getNextConstraintEmailValue(emailColumn); + emailColumnData.setValue(newEmail); + searchTextColumnData.setValue(searchTextColumnData.getNextConstraintEmailValue(searchTextColumn)); + String id = UUIDConverter.fromString(this.getColumnData(data, "id").getValue()).toString(); + log.warn("Found user with duplicate email [id:[{}]]. Attempting to rename email from '{}' to '{}'...", id, prevEmail, newEmail); + } + + protected void ignoreRecord(List batchData, CassandraToSqlColumnData[] data) { + log.warn("[{}] Affected data:\n{}", this.sqlTableName, this.dataToString(data)); + int index = batchData.indexOf(data); + if (index > 0) { + batchData.remove(index); + } + } + + protected CassandraToSqlColumn getColumn(String sqlColumnName) { + return this.columns.stream().filter(col -> col.getSqlColumnName().equals(sqlColumnName)).findFirst().get(); + } + + protected CassandraToSqlColumnData getColumnData(CassandraToSqlColumnData[] data, String sqlColumnName) { + CassandraToSqlColumn column = this.getColumn(sqlColumnName); + return data[column.getIndex()]; + } + + private Optional extractConstraintName(SQLException ex) { + final String sqlState = JdbcExceptionHelper.extractSqlState( ex ); + if (sqlState != null) { + String sqlStateClassCode = JdbcExceptionHelper.determineSqlStateClassCode( sqlState ); + if ( sqlStateClassCode != null ) { + if (Arrays.asList( + "23", // "integrity constraint violation" + "27", // "triggered data change violation" + "44" // "with check option violation" + ).contains(sqlStateClassCode)) { + if (ex instanceof PSQLException) { + return Optional.of(((PSQLException)ex).getServerErrorMessage().getConstraint()); + } + } + } + } + return Optional.empty(); } private Statement createCassandraSelectStatement() { @@ -103,8 +292,13 @@ public class CassandraToSqlTable { } insertStatementBuilder.deleteCharAt(insertStatementBuilder.length() - 1); insertStatementBuilder.append(") VALUES ("); - for (CassandraToSqlColumn ignored : columns) { - insertStatementBuilder.append("?").append(","); + for (CassandraToSqlColumn column : columns) { + if (column.getType() == CassandraToSqlColumnType.JSON) { + insertStatementBuilder.append("cast(? AS json)"); + } else { + insertStatementBuilder.append("?"); + } + insertStatementBuilder.append(","); } insertStatementBuilder.deleteCharAt(insertStatementBuilder.length() - 1); insertStatementBuilder.append(")");