This is an automated email from the ASF dual-hosted git repository.
azexin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere.git
The following commit(s) were added to refs/heads/master by this push:
new 55b641bf66a pipeline job support multiple columns unique key table
(#24161)
55b641bf66a is described below
commit 55b641bf66a6093feab8ec337898c6f7492a4ed6
Author: Hongsheng Zhong <[email protected]>
AuthorDate: Tue Feb 14 19:22:49 2023 +0800
pipeline job support multiple columns unique key table (#24161)
* Remove tableNameSchemaNameMapping from MigrationDataConsistencyChecker
* Load uniqueKeyColumn dynamically in MigrationDataConsistencyChecker;
Remove uniqueKeyColumn in migration job configuration;
* pipeline job support multiple columns unique key table
---
.../ingest/InventoryDumperConfiguration.java | 16 ++++-
.../api/config/job/MigrationJobConfiguration.java | 3 -
.../ingest/position/PrimaryKeyPositionFactory.java | 3 +-
.../spi/sqlbuilder/PipelineSQLBuilder.java | 26 ++++++---
.../core/ingest/dumper/InventoryDumper.java | 68 ++++++++++++++--------
.../core/ingest/exception/IngestException.java | 1 +
.../metadata/loader/PipelineTableMetaDataUtil.java | 33 +++++------
.../core/prepare/InventoryTaskSplitter.java | 51 +++++++++-------
.../sqlbuilder/AbstractPipelineSQLBuilder.java | 29 +++++----
.../yaml/job/YamlMigrationJobConfiguration.java | 3 -
.../job/YamlMigrationJobConfigurationSwapper.java | 6 +-
.../YamlPipelineColumnMetaDataSwapper.java | 52 -----------------
.../fixture/FixturePipelineSQLBuilder.java | 13 +++--
.../migration/api/impl/MigrationJobAPI.java | 10 ----
.../MigrationDataConsistencyChecker.java | 19 +++---
.../migration/prepare/MigrationJobPreparer.java | 4 --
.../primarykey/IndexesMigrationE2EIT.java | 26 +++++++++
.../primarykey/TextPrimaryKeyMigrationE2EIT.java | 1 -
.../env/scenario/primary_key/unique_key/mysql.xml | 28 ---------
.../api/impl/GovernanceRepositoryAPIImplTest.java | 5 +-
.../core/prepare/InventoryTaskSplitterTest.java | 55 ++++++++---------
.../data/pipeline/core/task/InventoryTaskTest.java | 5 +-
.../core/util/JobConfigurationBuilder.java | 14 -----
.../pipeline/core/util/PipelineContextUtil.java | 10 ++++
24 files changed, 224 insertions(+), 257 deletions(-)
diff --git
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ingest/InventoryDumperConfiguration.java
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ingest/InventoryDumperConfiguration.java
index 9340b0e7208..b52539ac4c8 100644
---
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ingest/InventoryDumperConfiguration.java
+++
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ingest/InventoryDumperConfiguration.java
@@ -20,8 +20,11 @@ package
org.apache.shardingsphere.data.pipeline.api.config.ingest;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
+import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
import
org.apache.shardingsphere.data.pipeline.spi.ratelimit.JobRateLimitAlgorithm;
+import java.util.List;
+
/**
* Inventory dumper configuration.
*/
@@ -35,9 +38,7 @@ public final class InventoryDumperConfiguration extends
DumperConfiguration {
private String logicTableName;
- private String uniqueKey;
-
- private Integer uniqueKeyDataType;
+ private List<PipelineColumnMetaData> uniqueKeyColumns;
private Integer shardingItem;
@@ -51,4 +52,13 @@ public final class InventoryDumperConfiguration extends
DumperConfiguration {
setTableNameMap(dumperConfig.getTableNameMap());
setTableNameSchemaNameMapping(dumperConfig.getTableNameSchemaNameMapping());
}
+
+ /**
+ * Has unique key or not.
+ *
+ * @return true when there's unique key, else false
+ */
+ public boolean hasUniqueKey() {
+ return null != uniqueKeyColumns && uniqueKeyColumns.size() > 0;
+ }
}
diff --git
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/job/MigrationJobConfiguration.java
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/job/MigrationJobConfiguration.java
index 15e6c0feeda..ddbd3478fe5 100644
---
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/job/MigrationJobConfiguration.java
+++
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/job/MigrationJobConfiguration.java
@@ -21,7 +21,6 @@ import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.ToString;
import
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
-import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
import java.util.List;
@@ -66,8 +65,6 @@ public final class MigrationJobConfiguration implements
PipelineJobConfiguration
private final List<String> jobShardingDataNodes;
- private final PipelineColumnMetaData uniqueKeyColumn;
-
private final int concurrency;
private final int retryTimes;
diff --git
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/ingest/position/PrimaryKeyPositionFactory.java
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/ingest/position/PrimaryKeyPositionFactory.java
index 06d20776135..cbbf3cd6faf 100644
---
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/ingest/position/PrimaryKeyPositionFactory.java
+++
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/ingest/position/PrimaryKeyPositionFactory.java
@@ -66,8 +66,9 @@ public final class PrimaryKeyPositionFactory {
return new IntegerPrimaryKeyPosition(((Number)
beginValue).longValue(), ((Number) endValue).longValue());
}
if (beginValue instanceof CharSequence) {
- return new StringPrimaryKeyPosition(beginValue.toString(),
endValue.toString());
+ return new StringPrimaryKeyPosition(beginValue.toString(), null !=
endValue ? endValue.toString() : null);
}
+ // TODO support more types, e.g. byte[] (MySQL varbinary)
return new UnsupportedKeyPosition();
}
}
diff --git
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/sqlbuilder/PipelineSQLBuilder.java
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/sqlbuilder/PipelineSQLBuilder.java
index edff4ca1d18..9c042a8bd35 100644
---
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/sqlbuilder/PipelineSQLBuilder.java
+++
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/sqlbuilder/PipelineSQLBuilder.java
@@ -41,15 +41,24 @@ public interface PipelineSQLBuilder extends TypedSPI {
}
/**
- * Build divisible inventory dump first SQL.
+ * Build divisible inventory dump SQL.
*
* @param schemaName schema name
* @param tableName table name
* @param uniqueKey unique key
- * @param uniqueKeyDataType unique key data type
* @return divisible inventory dump SQL
*/
- String buildDivisibleInventoryDumpSQL(String schemaName, String tableName,
String uniqueKey, int uniqueKeyDataType);
+ String buildDivisibleInventoryDumpSQL(String schemaName, String tableName,
String uniqueKey);
+
+ /**
+ * Build divisible inventory dump SQL without end value.
+ *
+ * @param schemaName schema name
+ * @param tableName table name
+ * @param uniqueKey unique key
+ * @return divisible inventory dump SQL without end value
+ */
+ String buildDivisibleInventoryDumpSQLNoEnd(String schemaName, String
tableName, String uniqueKey);
/**
* Build indivisible inventory dump first SQL.
@@ -57,19 +66,18 @@ public interface PipelineSQLBuilder extends TypedSPI {
* @param schemaName schema name
* @param tableName table name
* @param uniqueKey unique key
- * @param uniqueKeyDataType unique key data type
* @return indivisible inventory dump SQL
*/
- String buildIndivisibleInventoryDumpSQL(String schemaName, String
tableName, String uniqueKey, int uniqueKeyDataType);
+ String buildIndivisibleInventoryDumpSQL(String schemaName, String
tableName, String uniqueKey);
/**
- * Build inventory dump all SQL.
+ * Build no unique key inventory dump SQL.
*
* @param schemaName schema name
* @param tableName tableName
* @return inventory dump all SQL
*/
- String buildInventoryDumpAllSQL(String schemaName, String tableName);
+ String buildNoUniqueKeyInventoryDumpSQL(String schemaName, String
tableName);
/**
* Build insert SQL.
@@ -152,10 +160,10 @@ public interface PipelineSQLBuilder extends TypedSPI {
*
* @param schemaName schema name
* @param tableName table name
- * @param primaryKey primary key
+ * @param uniqueKey unique key
* @return split SQL
*/
- String buildSplitByPrimaryKeyRangeSQL(String schemaName, String tableName,
String primaryKey);
+ String buildSplitByPrimaryKeyRangeSQL(String schemaName, String tableName,
String uniqueKey);
/**
* Build CRC32 SQL.
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/dumper/InventoryDumper.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/dumper/InventoryDumper.java
index 5b5a1642c93..08f43c64f5b 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/dumper/InventoryDumper.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/dumper/InventoryDumper.java
@@ -36,6 +36,7 @@ import
org.apache.shardingsphere.data.pipeline.api.ingest.record.FinishedRecord;
import org.apache.shardingsphere.data.pipeline.api.job.JobOperationType;
import org.apache.shardingsphere.data.pipeline.api.metadata.LogicTableName;
import
org.apache.shardingsphere.data.pipeline.api.metadata.loader.PipelineTableMetaDataLoader;
+import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineTableMetaData;
import
org.apache.shardingsphere.data.pipeline.core.ingest.IngestDataChangeType;
import
org.apache.shardingsphere.data.pipeline.core.ingest.exception.IngestException;
@@ -98,39 +99,28 @@ public final class InventoryDumper extends
AbstractLifecycleExecutor implements
}
PipelineTableMetaData tableMetaData =
metaDataLoader.getTableMetaData(dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName())),
dumperConfig.getActualTableName());
try (Connection connection = dataSource.getConnection()) {
- dump(tableMetaData, connection, buildInventoryDumpSQL(),
((PrimaryKeyPosition<?>) position).getBeginValue());
+ dump(tableMetaData, connection);
log.info("Inventory dump done");
} catch (final SQLException ex) {
log.error("Inventory dump, ex caught, msg={}.", ex.getMessage());
- throw new IngestException(ex);
+ throw new IngestException("Inventory dump failed on " +
dumperConfig.getActualTableName(), ex);
} finally {
channel.pushRecord(new FinishedRecord(new FinishedPosition()));
}
}
- private String buildInventoryDumpSQL() {
- String schemaName = dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName()));
- if (null == dumperConfig.getUniqueKey()) {
- return sqlBuilder.buildInventoryDumpAllSQL(schemaName,
dumperConfig.getActualTableName());
- }
- if
(PipelineJdbcUtils.isIntegerColumn(dumperConfig.getUniqueKeyDataType())) {
- return sqlBuilder.buildDivisibleInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName(), dumperConfig.getUniqueKey(),
dumperConfig.getUniqueKeyDataType());
- }
- return sqlBuilder.buildIndivisibleInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName(), dumperConfig.getUniqueKey(),
dumperConfig.getUniqueKeyDataType());
- }
-
- private void dump(final PipelineTableMetaData tableMetaData, final
Connection connection, final String sql, final Object beginUniqueKeyValue)
throws SQLException {
+ private void dump(final PipelineTableMetaData tableMetaData, final
Connection connection) throws SQLException {
if (null != dumperConfig.getRateLimitAlgorithm()) {
dumperConfig.getRateLimitAlgorithm().intercept(JobOperationType.SELECT, 1);
}
int batchSize = dumperConfig.getBatchSize();
DatabaseType databaseType =
dumperConfig.getDataSourceConfig().getDatabaseType();
- try (PreparedStatement preparedStatement =
JDBCStreamQueryUtil.generateStreamQueryPreparedStatement(databaseType,
connection, sql)) {
+ try (PreparedStatement preparedStatement =
JDBCStreamQueryUtil.generateStreamQueryPreparedStatement(databaseType,
connection, buildInventoryDumpSQL())) {
dumpStatement = preparedStatement;
if (!(databaseType instanceof MySQLDatabaseType)) {
preparedStatement.setFetchSize(batchSize);
}
- setParameters(preparedStatement, beginUniqueKeyValue);
+ setParameters(preparedStatement);
try (ResultSet resultSet = preparedStatement.executeQuery()) {
ResultSetMetaData resultSetMetaData = resultSet.getMetaData();
while (resultSet.next()) {
@@ -145,13 +135,45 @@ public final class InventoryDumper extends
AbstractLifecycleExecutor implements
}
}
- private void setParameters(final PreparedStatement preparedStatement,
final Object beginUniqueKeyValue) throws SQLException {
- if (null == dumperConfig.getUniqueKey()) {
+ private String buildInventoryDumpSQL() {
+ String schemaName = dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName()));
+ if (!dumperConfig.hasUniqueKey()) {
+ return sqlBuilder.buildNoUniqueKeyInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName());
+ }
+ PipelineColumnMetaData firstColumn =
dumperConfig.getUniqueKeyColumns().get(0);
+ if (PipelineJdbcUtils.isIntegerColumn(firstColumn.getDataType())) {
+ return sqlBuilder.buildDivisibleInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName(), firstColumn.getName());
+ }
+ if (PipelineJdbcUtils.isStringColumn(firstColumn.getDataType())) {
+ PrimaryKeyPosition<?> position = (PrimaryKeyPosition<?>)
dumperConfig.getPosition();
+ if (null != position.getBeginValue() && null !=
position.getEndValue()) {
+ return sqlBuilder.buildDivisibleInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName(), firstColumn.getName());
+ }
+ if (null != position.getBeginValue() && null ==
position.getEndValue()) {
+ return
sqlBuilder.buildDivisibleInventoryDumpSQLNoEnd(schemaName,
dumperConfig.getActualTableName(), firstColumn.getName());
+ }
+ }
+ return sqlBuilder.buildIndivisibleInventoryDumpSQL(schemaName,
dumperConfig.getActualTableName(), firstColumn.getName());
+ }
+
+ private void setParameters(final PreparedStatement preparedStatement)
throws SQLException {
+ if (!dumperConfig.hasUniqueKey()) {
return;
}
- if
(PipelineJdbcUtils.isIntegerColumn(dumperConfig.getUniqueKeyDataType())) {
- preparedStatement.setObject(1, beginUniqueKeyValue);
- preparedStatement.setObject(2, ((PrimaryKeyPosition<?>)
dumperConfig.getPosition()).getEndValue());
+ PipelineColumnMetaData firstColumn =
dumperConfig.getUniqueKeyColumns().get(0);
+ PrimaryKeyPosition<?> position = (PrimaryKeyPosition<?>)
dumperConfig.getPosition();
+ if (PipelineJdbcUtils.isIntegerColumn(firstColumn.getDataType())) {
+ preparedStatement.setObject(1, position.getBeginValue());
+ preparedStatement.setObject(2, position.getEndValue());
+ return;
+ }
+ if (PipelineJdbcUtils.isStringColumn(firstColumn.getDataType())) {
+ if (null != position.getBeginValue()) {
+ preparedStatement.setObject(1, position.getBeginValue());
+ }
+ if (null != position.getEndValue()) {
+ preparedStatement.setObject(2, position.getEndValue());
+ }
}
}
@@ -167,9 +189,9 @@ public final class InventoryDumper extends
AbstractLifecycleExecutor implements
}
private IngestPosition<?> newPosition(final ResultSet resultSet) throws
SQLException {
- return null == dumperConfig.getUniqueKey()
+ return !dumperConfig.hasUniqueKey()
? new PlaceholderPosition()
- :
PrimaryKeyPositionFactory.newInstance(resultSet.getObject(dumperConfig.getUniqueKey()),
((PrimaryKeyPosition<?>) dumperConfig.getPosition()).getEndValue());
+ :
PrimaryKeyPositionFactory.newInstance(resultSet.getObject(dumperConfig.getUniqueKeyColumns().get(0).getName()),
((PrimaryKeyPosition<?>) dumperConfig.getPosition()).getEndValue());
}
@Override
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/exception/IngestException.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/exception/IngestException.java
index 527c4db6ef8..317ffe5750b 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/exception/IngestException.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/ingest/exception/IngestException.java
@@ -20,6 +20,7 @@ package
org.apache.shardingsphere.data.pipeline.core.ingest.exception;
/**
* Ingest exception.
*/
+// TODO extends from PipelineSQLException
public final class IngestException extends RuntimeException {
private static final long serialVersionUID = 1L;
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/metadata/loader/PipelineTableMetaDataUtil.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/metadata/loader/PipelineTableMetaDataUtil.java
index ff652e93280..84482793b4b 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/metadata/loader/PipelineTableMetaDataUtil.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/metadata/loader/PipelineTableMetaDataUtil.java
@@ -27,8 +27,9 @@ import
org.apache.shardingsphere.data.pipeline.core.exception.job.SplitPipelineJ
import
org.apache.shardingsphere.infra.util.exception.ShardingSpherePreconditions;
import java.util.Collection;
+import java.util.LinkedList;
import java.util.List;
-import java.util.Optional;
+import java.util.stream.Collectors;
/**
* Pipeline table meta data util.
@@ -37,34 +38,30 @@ import java.util.Optional;
public final class PipelineTableMetaDataUtil {
/**
- * Get unique key column.
+ * Get unique key columns.
*
* @param schemaName schema name
* @param tableName table name
* @param metaDataLoader meta data loader
- * @return pipeline column meta data
+ * @return pipeline columns meta data
*/
- public static Optional<PipelineColumnMetaData> getUniqueKeyColumn(final
String schemaName, final String tableName, final PipelineTableMetaDataLoader
metaDataLoader) {
- PipelineTableMetaData pipelineTableMetaData =
metaDataLoader.getTableMetaData(schemaName, tableName);
- return
Optional.ofNullable(getAnAppropriateUniqueKeyColumn(pipelineTableMetaData,
tableName));
- }
-
- private static PipelineColumnMetaData
getAnAppropriateUniqueKeyColumn(final PipelineTableMetaData tableMetaData,
final String tableName) {
+ public static List<PipelineColumnMetaData> getUniqueKeyColumns(final
String schemaName, final String tableName, final PipelineTableMetaDataLoader
metaDataLoader) {
+ PipelineTableMetaData tableMetaData =
metaDataLoader.getTableMetaData(schemaName, tableName);
ShardingSpherePreconditions.checkNotNull(tableMetaData, () -> new
SplitPipelineJobByRangeException(tableName, "Can not get table meta data"));
List<String> primaryKeys = tableMetaData.getPrimaryKeyColumns();
- if (1 == primaryKeys.size()) {
- return
tableMetaData.getColumnMetaData(tableMetaData.getPrimaryKeyColumns().get(0));
+ if (primaryKeys.size() > 0) {
+ return
primaryKeys.stream().map(tableMetaData::getColumnMetaData).collect(Collectors.toList());
}
Collection<PipelineIndexMetaData> uniqueIndexes =
tableMetaData.getUniqueIndexes();
- if (uniqueIndexes.isEmpty() && primaryKeys.isEmpty()) {
- return null;
+ if (uniqueIndexes.isEmpty()) {
+ return new LinkedList<>();
}
- if (1 == uniqueIndexes.size() && 1 ==
uniqueIndexes.iterator().next().getColumns().size()) {
- PipelineColumnMetaData column =
uniqueIndexes.iterator().next().getColumns().get(0);
- if (!column.isNullable()) {
- return column;
+ for (PipelineIndexMetaData each : uniqueIndexes) {
+ if
(each.getColumns().stream().anyMatch(PipelineColumnMetaData::isNullable)) {
+ continue;
}
+ return each.getColumns();
}
- throw new SplitPipelineJobByRangeException(tableName, "table contains
multiple unique index or unique index contains nullable/multiple column(s)");
+ return new LinkedList<>();
}
}
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/prepare/InventoryTaskSplitter.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/prepare/InventoryTaskSplitter.java
index 49c2e98ac5c..b1c8433e99e 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/prepare/InventoryTaskSplitter.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/prepare/InventoryTaskSplitter.java
@@ -17,7 +17,6 @@
package org.apache.shardingsphere.data.pipeline.core.prepare;
-import com.google.common.base.Strings;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import
org.apache.shardingsphere.data.pipeline.api.config.ImporterConfiguration;
@@ -30,6 +29,7 @@ import
org.apache.shardingsphere.data.pipeline.api.ingest.position.IntegerPrimar
import
org.apache.shardingsphere.data.pipeline.api.ingest.position.NoUniqueKeyPosition;
import
org.apache.shardingsphere.data.pipeline.api.ingest.position.PlaceholderPosition;
import
org.apache.shardingsphere.data.pipeline.api.ingest.position.StringPrimaryKeyPosition;
+import
org.apache.shardingsphere.data.pipeline.api.ingest.position.UnsupportedKeyPosition;
import org.apache.shardingsphere.data.pipeline.api.job.JobStatus;
import
org.apache.shardingsphere.data.pipeline.api.job.progress.InventoryIncrementalJobItemProgress;
import org.apache.shardingsphere.data.pipeline.api.metadata.LogicTableName;
@@ -54,7 +54,6 @@ import java.util.Collection;
import java.util.Collections;
import java.util.LinkedList;
import java.util.List;
-import java.util.Optional;
/**
* Inventory data task splitter.
@@ -103,8 +102,7 @@ public final class InventoryTaskSplitter {
inventoryDumperConfig.setActualTableName(key.getOriginal());
inventoryDumperConfig.setLogicTableName(value.getOriginal());
inventoryDumperConfig.setPosition(new PlaceholderPosition());
- inventoryDumperConfig.setUniqueKey(dumperConfig.getUniqueKey());
-
inventoryDumperConfig.setUniqueKeyDataType(dumperConfig.getUniqueKeyDataType());
+
inventoryDumperConfig.setUniqueKeyColumns(dumperConfig.getUniqueKeyColumns());
result.add(inventoryDumperConfig);
});
return result;
@@ -112,14 +110,11 @@ public final class InventoryTaskSplitter {
private Collection<InventoryDumperConfiguration> splitByPrimaryKey(final
InventoryDumperConfiguration dumperConfig, final
InventoryIncrementalJobItemContext jobItemContext,
final
DataSource dataSource) {
- if (null == dumperConfig.getUniqueKey()) {
+ if (null == dumperConfig.getUniqueKeyColumns()) {
String schemaName = dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName()));
String actualTableName = dumperConfig.getActualTableName();
- Optional<PipelineColumnMetaData> uniqueKeyColumn =
PipelineTableMetaDataUtil.getUniqueKeyColumn(schemaName, actualTableName,
jobItemContext.getSourceMetaDataLoader());
- uniqueKeyColumn.ifPresent(column -> {
- dumperConfig.setUniqueKey(column.getName());
- dumperConfig.setUniqueKeyDataType(column.getDataType());
- });
+ List<PipelineColumnMetaData> uniqueKeyColumns =
PipelineTableMetaDataUtil.getUniqueKeyColumns(schemaName, actualTableName,
jobItemContext.getSourceMetaDataLoader());
+ dumperConfig.setUniqueKeyColumns(uniqueKeyColumns);
}
Collection<InventoryDumperConfiguration> result = new LinkedList<>();
InventoryIncrementalProcessContext jobProcessContext =
jobItemContext.getJobProcessContext();
@@ -134,8 +129,7 @@ public final class InventoryTaskSplitter {
splitDumperConfig.setShardingItem(i++);
splitDumperConfig.setActualTableName(dumperConfig.getActualTableName());
splitDumperConfig.setLogicTableName(dumperConfig.getLogicTableName());
- splitDumperConfig.setUniqueKey(dumperConfig.getUniqueKey());
-
splitDumperConfig.setUniqueKeyDataType(dumperConfig.getUniqueKeyDataType());
+
splitDumperConfig.setUniqueKeyColumns(dumperConfig.getUniqueKeyColumns());
splitDumperConfig.setBatchSize(batchSize);
splitDumperConfig.setRateLimitAlgorithm(rateLimitAlgorithm);
result.add(splitDumperConfig);
@@ -150,14 +144,20 @@ public final class InventoryTaskSplitter {
// Do NOT filter FinishedPosition here, since whole inventory
tasks are required in job progress when persisting to register center.
return
initProgress.getInventory().getInventoryPosition(dumperConfig.getActualTableName()).values();
}
- if (Strings.isNullOrEmpty(dumperConfig.getUniqueKey())) {
+ if (!dumperConfig.hasUniqueKey()) {
return getPositionWithoutUniqueKey(jobItemContext, dataSource,
dumperConfig);
}
- int uniqueKeyDataType = dumperConfig.getUniqueKeyDataType();
- if (PipelineJdbcUtils.isIntegerColumn(uniqueKeyDataType)) {
- return getPositionByIntegerUniqueKeyRange(jobItemContext,
dataSource, dumperConfig);
+ List<PipelineColumnMetaData> uniqueKeyColumns =
dumperConfig.getUniqueKeyColumns();
+ if (1 == uniqueKeyColumns.size()) {
+ int firstColumnDataType = uniqueKeyColumns.get(0).getDataType();
+ if (PipelineJdbcUtils.isIntegerColumn(firstColumnDataType)) {
+ return getPositionByIntegerUniqueKeyRange(jobItemContext,
dataSource, dumperConfig);
+ }
+ if (PipelineJdbcUtils.isStringColumn(firstColumnDataType)) {
+ return getPositionByStringUniqueKeyRange(jobItemContext,
dataSource, dumperConfig);
+ }
}
- return getPositionByStringUniqueKeyRange(jobItemContext, dataSource,
dumperConfig);
+ return getUnsupportedPosition(jobItemContext, dataSource,
dumperConfig);
}
private Collection<IngestPosition<?>> getPositionWithoutUniqueKey(final
InventoryIncrementalJobItemContext jobItemContext, final DataSource dataSource,
@@ -181,7 +181,8 @@ public final class InventoryTaskSplitter {
return resultSet.getLong(1);
}
} catch (final SQLException ex) {
- throw new
SplitPipelineJobByUniqueKeyException(dumperConfig.getActualTableName(),
dumperConfig.getUniqueKey(), ex);
+ String uniqueKey = dumperConfig.hasUniqueKey() ?
dumperConfig.getUniqueKeyColumns().get(0).getName() : "";
+ throw new
SplitPipelineJobByUniqueKeyException(dumperConfig.getActualTableName(),
uniqueKey, ex);
}
}
@@ -189,8 +190,9 @@ public final class InventoryTaskSplitter {
final InventoryDumperConfiguration dumperConfig) {
Collection<IngestPosition<?>> result = new LinkedList<>();
PipelineJobConfiguration jobConfig = jobItemContext.getJobConfig();
+ String uniqueKey = dumperConfig.getUniqueKeyColumns().get(0).getName();
String sql = TypedSPILoader.getService(PipelineSQLBuilder.class,
jobConfig.getSourceDatabaseType())
- .buildSplitByPrimaryKeyRangeSQL(dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName())),
dumperConfig.getActualTableName(), dumperConfig.getUniqueKey());
+ .buildSplitByPrimaryKeyRangeSQL(dumperConfig.getSchemaName(new
LogicTableName(dumperConfig.getLogicTableName())),
dumperConfig.getActualTableName(), uniqueKey);
int shardingSize =
jobItemContext.getJobProcessContext().getPipelineProcessConfig().getRead().getShardingSize();
try (
Connection connection = dataSource.getConnection();
@@ -220,7 +222,7 @@ public final class InventoryTaskSplitter {
result.add(new IntegerPrimaryKeyPosition(0, 0));
}
} catch (final SQLException ex) {
- throw new
SplitPipelineJobByUniqueKeyException(dumperConfig.getActualTableName(),
dumperConfig.getUniqueKey(), ex);
+ throw new
SplitPipelineJobByUniqueKeyException(dumperConfig.getActualTableName(),
uniqueKey, ex);
}
return result;
}
@@ -230,7 +232,14 @@ public final class InventoryTaskSplitter {
long tableRecordsCount = getTableRecordsCount(jobItemContext,
dataSource, dumperConfig);
jobItemContext.updateInventoryRecordsCount(tableRecordsCount);
Collection<IngestPosition<?>> result = new LinkedList<>();
- result.add(new StringPrimaryKeyPosition("!", "~"));
+ result.add(new StringPrimaryKeyPosition(null, null));
return result;
}
+
+ private Collection<IngestPosition<?>> getUnsupportedPosition(final
InventoryIncrementalJobItemContext jobItemContext, final DataSource dataSource,
+ final
InventoryDumperConfiguration dumperConfig) {
+ long tableRecordsCount = getTableRecordsCount(jobItemContext,
dataSource, dumperConfig);
+ jobItemContext.updateInventoryRecordsCount(tableRecordsCount);
+ return Collections.singletonList(new UnsupportedKeyPosition());
+ }
}
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/sqlbuilder/AbstractPipelineSQLBuilder.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/sqlbuilder/AbstractPipelineSQLBuilder.java
index ac3455e9b74..f6290fc4fc9 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/sqlbuilder/AbstractPipelineSQLBuilder.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/sqlbuilder/AbstractPipelineSQLBuilder.java
@@ -71,19 +71,32 @@ public abstract class AbstractPipelineSQLBuilder implements
PipelineSQLBuilder {
protected abstract String getRightIdentifierQuoteString();
@Override
- public String buildDivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey, final int uniqueKeyDataType) {
+ public String buildDivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey) {
String qualifiedTableName = getQualifiedTableName(schemaName,
tableName);
String quotedUniqueKey = quote(uniqueKey);
return String.format("SELECT * FROM %s WHERE %s>=? AND %s<=? ORDER BY
%s ASC", qualifiedTableName, quotedUniqueKey, quotedUniqueKey, quotedUniqueKey);
}
@Override
- public String buildIndivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey, final int uniqueKeyDataType) {
+ public String buildDivisibleInventoryDumpSQLNoEnd(final String schemaName,
final String tableName, final String uniqueKey) {
+ String qualifiedTableName = getQualifiedTableName(schemaName,
tableName);
+ String quotedUniqueKey = quote(uniqueKey);
+ return String.format("SELECT * FROM %s WHERE %s>=? ORDER BY %s ASC",
qualifiedTableName, quotedUniqueKey, quotedUniqueKey);
+ }
+
+ @Override
+ public String buildIndivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey) {
String qualifiedTableName = getQualifiedTableName(schemaName,
tableName);
String quotedUniqueKey = quote(uniqueKey);
return String.format("SELECT * FROM %s ORDER BY %s ASC",
qualifiedTableName, quotedUniqueKey);
}
+ @Override
+ public String buildNoUniqueKeyInventoryDumpSQL(final String schemaName,
final String tableName) {
+ String qualifiedTableName = getQualifiedTableName(schemaName,
tableName);
+ return String.format("SELECT * FROM %s", qualifiedTableName);
+ }
+
protected final String getQualifiedTableName(final String schemaName,
final String tableName) {
StringBuilder result = new StringBuilder();
if (TypedSPILoader.getService(DatabaseType.class,
getType()).isSchemaAvailable() && !Strings.isNullOrEmpty(schemaName)) {
@@ -184,15 +197,9 @@ public abstract class AbstractPipelineSQLBuilder
implements PipelineSQLBuilder {
}
@Override
- public String buildSplitByPrimaryKeyRangeSQL(final String schemaName,
final String tableName, final String primaryKey) {
- String quotedUniqueKey = quote(primaryKey);
- return String.format("SELECT MAX(%s),COUNT(*) FROM (SELECT %s FROM %s
WHERE %s>=? ORDER BY %s LIMIT ?) t",
+ public String buildSplitByPrimaryKeyRangeSQL(final String schemaName,
final String tableName, final String uniqueKey) {
+ String quotedUniqueKey = quote(uniqueKey);
+ return String.format("SELECT MAX(%s),COUNT(1) FROM (SELECT %s FROM %s
WHERE %s>=? ORDER BY %s LIMIT ?) t",
quotedUniqueKey, quotedUniqueKey,
getQualifiedTableName(schemaName, tableName), quotedUniqueKey, quotedUniqueKey);
}
-
- @Override
- public String buildInventoryDumpAllSQL(final String schemaName, final
String tableName) {
- String qualifiedTableName = getQualifiedTableName(schemaName,
tableName);
- return String.format("SELECT * FROM %s", qualifiedTableName);
- }
}
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
index e5a8763548a..eca1818a946 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
@@ -23,7 +23,6 @@ import lombok.Setter;
import lombok.ToString;
import
org.apache.shardingsphere.data.pipeline.api.config.job.yaml.YamlPipelineJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.api.datasource.config.yaml.YamlPipelineDataSourceConfiguration;
-import
org.apache.shardingsphere.data.pipeline.yaml.metadata.YamlPipelineColumnMetaData;
import java.util.List;
@@ -66,8 +65,6 @@ public final class YamlMigrationJobConfiguration implements
YamlPipelineJobConfi
private List<String> jobShardingDataNodes;
- private YamlPipelineColumnMetaData uniqueKeyColumn;
-
private int concurrency = 3;
private int retryTimes = 3;
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
index a88359b552c..601c300abeb 100644
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
+++
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
@@ -19,7 +19,6 @@ package org.apache.shardingsphere.data.pipeline.yaml.job;
import
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.api.datasource.config.yaml.YamlPipelineDataSourceConfigurationSwapper;
-import
org.apache.shardingsphere.data.pipeline.yaml.metadata.YamlPipelineColumnMetaDataSwapper;
import org.apache.shardingsphere.infra.util.yaml.YamlEngine;
import
org.apache.shardingsphere.infra.util.yaml.swapper.YamlConfigurationSwapper;
@@ -30,8 +29,6 @@ public final class YamlMigrationJobConfigurationSwapper
implements YamlConfigura
private final YamlPipelineDataSourceConfigurationSwapper
dataSourceConfigSwapper = new YamlPipelineDataSourceConfigurationSwapper();
- private final YamlPipelineColumnMetaDataSwapper
pipelineColumnMetaDataSwapper = new YamlPipelineColumnMetaDataSwapper();
-
@Override
public YamlMigrationJobConfiguration swapToYamlConfiguration(final
MigrationJobConfiguration data) {
YamlMigrationJobConfiguration result = new
YamlMigrationJobConfiguration();
@@ -47,7 +44,6 @@ public final class YamlMigrationJobConfigurationSwapper
implements YamlConfigura
result.setTarget(dataSourceConfigSwapper.swapToYamlConfiguration(data.getTarget()));
result.setTablesFirstDataNodes(data.getTablesFirstDataNodes());
result.setJobShardingDataNodes(data.getJobShardingDataNodes());
-
result.setUniqueKeyColumn(pipelineColumnMetaDataSwapper.swapToYamlConfiguration(data.getUniqueKeyColumn()));
result.setConcurrency(data.getConcurrency());
result.setRetryTimes(data.getRetryTimes());
return result;
@@ -60,7 +56,7 @@ public final class YamlMigrationJobConfigurationSwapper
implements YamlConfigura
yamlConfig.getSourceDatabaseType(),
yamlConfig.getTargetDatabaseType(),
yamlConfig.getSourceTableName(),
yamlConfig.getTargetTableName(),
dataSourceConfigSwapper.swapToObject(yamlConfig.getSource()),
dataSourceConfigSwapper.swapToObject(yamlConfig.getTarget()),
- yamlConfig.getTablesFirstDataNodes(),
yamlConfig.getJobShardingDataNodes(),
pipelineColumnMetaDataSwapper.swapToObject(yamlConfig.getUniqueKeyColumn()),
+ yamlConfig.getTablesFirstDataNodes(),
yamlConfig.getJobShardingDataNodes(),
yamlConfig.getConcurrency(), yamlConfig.getRetryTimes());
}
diff --git
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/metadata/YamlPipelineColumnMetaDataSwapper.java
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/metadata/YamlPipelineColumnMetaDataSwapper.java
deleted file mode 100644
index 645d20053d2..00000000000
---
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/metadata/YamlPipelineColumnMetaDataSwapper.java
+++ /dev/null
@@ -1,52 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.shardingsphere.data.pipeline.yaml.metadata;
-
-import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
-import
org.apache.shardingsphere.infra.util.yaml.swapper.YamlConfigurationSwapper;
-
-/**
- * Yaml pipeline column meta data swapper.
- */
-public final class YamlPipelineColumnMetaDataSwapper implements
YamlConfigurationSwapper<YamlPipelineColumnMetaData, PipelineColumnMetaData> {
-
- @Override
- public YamlPipelineColumnMetaData swapToYamlConfiguration(final
PipelineColumnMetaData data) {
- if (null == data) {
- return null;
- }
- YamlPipelineColumnMetaData result = new YamlPipelineColumnMetaData();
- result.setName(data.getName());
- result.setDataType(data.getDataType());
- result.setDataTypeName(data.getDataTypeName());
- result.setNullable(data.isNullable());
- result.setPrimaryKey(data.isPrimaryKey());
- result.setOrdinalPosition(data.getOrdinalPosition());
- result.setUniqueKey(data.isUniqueKey());
- return result;
- }
-
- @Override
- public PipelineColumnMetaData swapToObject(final
YamlPipelineColumnMetaData yamlConfig) {
- if (null == yamlConfig) {
- return null;
- }
- return new PipelineColumnMetaData(yamlConfig.getOrdinalPosition(),
yamlConfig.getName(), yamlConfig.getDataType(), yamlConfig.getDataTypeName(),
yamlConfig.isNullable(),
- yamlConfig.isPrimaryKey(), yamlConfig.isUniqueKey());
- }
-}
diff --git
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/core/check/consistency/algorithm/fixture/FixturePipelineSQLBuilder.java
b/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/core/check/consistency/algorithm/fixture/FixturePipelineSQLBuilder.java
index 08e46bfa21e..136723e894b 100644
---
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/core/check/consistency/algorithm/fixture/FixturePipelineSQLBuilder.java
+++
b/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/core/check/consistency/algorithm/fixture/FixturePipelineSQLBuilder.java
@@ -29,12 +29,17 @@ import java.util.Optional;
public final class FixturePipelineSQLBuilder implements PipelineSQLBuilder {
@Override
- public String buildDivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey, final int uniqueKeyDataType) {
+ public String buildDivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey) {
return "";
}
@Override
- public String buildIndivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey, final int uniqueKeyDataType) {
+ public String buildDivisibleInventoryDumpSQLNoEnd(final String schemaName,
final String tableName, final String uniqueKey) {
+ return "";
+ }
+
+ @Override
+ public String buildIndivisibleInventoryDumpSQL(final String schemaName,
final String tableName, final String uniqueKey) {
return "";
}
@@ -79,7 +84,7 @@ public final class FixturePipelineSQLBuilder implements
PipelineSQLBuilder {
}
@Override
- public String buildSplitByPrimaryKeyRangeSQL(final String schemaName,
final String tableName, final String primaryKey) {
+ public String buildSplitByPrimaryKeyRangeSQL(final String schemaName,
final String tableName, final String uniqueKey) {
return "";
}
@@ -89,7 +94,7 @@ public final class FixturePipelineSQLBuilder implements
PipelineSQLBuilder {
}
@Override
- public String buildInventoryDumpAllSQL(final String schemaName, final
String tableName) {
+ public String buildNoUniqueKeyInventoryDumpSQL(final String schemaName,
final String tableName) {
return "";
}
diff --git
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/api/impl/MigrationJobAPI.java
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/api/impl/MigrationJobAPI.java
index 02f479a4ae3..7f914dee3c8 100644
---
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/api/impl/MigrationJobAPI.java
+++
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/api/impl/MigrationJobAPI.java
@@ -60,8 +60,6 @@ import
org.apache.shardingsphere.data.pipeline.core.datasource.PipelineDataSourc
import
org.apache.shardingsphere.data.pipeline.core.exception.connection.RegisterMigrationSourceStorageUnitException;
import
org.apache.shardingsphere.data.pipeline.core.exception.connection.UnregisterMigrationSourceStorageUnitException;
import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.PipelineSchemaUtil;
-import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.PipelineTableMetaDataUtil;
-import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.StandardPipelineTableMetaDataLoader;
import
org.apache.shardingsphere.data.pipeline.core.sharding.ShardingColumnsExtractor;
import org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJob;
import
org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJobId;
@@ -74,7 +72,6 @@ import
org.apache.shardingsphere.data.pipeline.spi.job.JobTypeFactory;
import
org.apache.shardingsphere.data.pipeline.spi.sqlbuilder.PipelineSQLBuilder;
import
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigurationSwapper;
-import
org.apache.shardingsphere.data.pipeline.yaml.metadata.YamlPipelineColumnMetaDataSwapper;
import org.apache.shardingsphere.elasticjob.infra.pojo.JobConfigurationPOJO;
import org.apache.shardingsphere.infra.config.rule.RuleConfiguration;
import org.apache.shardingsphere.infra.database.metadata.DataSourceMetaData;
@@ -443,13 +440,6 @@ public final class MigrationJobAPI extends
AbstractInventoryIncrementalJobAPIImp
result.setTargetDatabaseType(targetPipelineDataSource.getDatabaseType().getType());
result.setTargetDatabaseName(targetDatabaseName);
result.setTargetTableName(param.getTargetTableName());
- try (PipelineDataSourceWrapper dataSource =
PipelineDataSourceFactory.newInstance(sourceDataSourceConfig)) {
- StandardPipelineTableMetaDataLoader metaDataLoader = new
StandardPipelineTableMetaDataLoader(dataSource);
- PipelineTableMetaDataUtil.getUniqueKeyColumn(sourceSchemaName,
param.getSourceTableName(), metaDataLoader)
- .ifPresent(optional -> result.setUniqueKeyColumn(new
YamlPipelineColumnMetaDataSwapper().swapToYamlConfiguration(optional)));
- } catch (final SQLException ex) {
- throw new RuntimeException(ex);
- }
extendYamlJobConfiguration(result);
MigrationJobConfiguration jobConfig = new
YamlMigrationJobConfigurationSwapper().swapToObject(result);
start(jobConfig);
diff --git
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyChecker.java
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyChecker.java
index a62c9cb83e9..d6254971c18 100644
---
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyChecker.java
+++
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyChecker.java
@@ -20,7 +20,6 @@ package
org.apache.shardingsphere.data.pipeline.scenario.migration.check.consist
import lombok.extern.slf4j.Slf4j;
import
org.apache.shardingsphere.data.pipeline.api.check.consistency.DataConsistencyCheckResult;
import
org.apache.shardingsphere.data.pipeline.api.check.consistency.PipelineDataConsistencyChecker;
-import
org.apache.shardingsphere.data.pipeline.api.config.TableNameSchemaNameMapping;
import
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
import
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
@@ -29,11 +28,13 @@ import
org.apache.shardingsphere.data.pipeline.api.metadata.SchemaName;
import org.apache.shardingsphere.data.pipeline.api.metadata.SchemaTableName;
import org.apache.shardingsphere.data.pipeline.api.metadata.TableName;
import
org.apache.shardingsphere.data.pipeline.api.metadata.loader.PipelineTableMetaDataLoader;
+import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
import
org.apache.shardingsphere.data.pipeline.core.check.consistency.ConsistencyCheckJobItemProgressContext;
import
org.apache.shardingsphere.data.pipeline.core.check.consistency.SingleTableInventoryDataConsistencyChecker;
import
org.apache.shardingsphere.data.pipeline.core.context.InventoryIncrementalProcessContext;
import
org.apache.shardingsphere.data.pipeline.core.datasource.PipelineDataSourceFactory;
import
org.apache.shardingsphere.data.pipeline.core.exception.data.UnsupportedPipelineDatabaseTypeException;
+import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.PipelineTableMetaDataUtil;
import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.StandardPipelineTableMetaDataLoader;
import
org.apache.shardingsphere.data.pipeline.scenario.migration.api.impl.MigrationJobAPI;
import
org.apache.shardingsphere.data.pipeline.spi.check.consistency.DataConsistencyCalculateAlgorithm;
@@ -42,9 +43,8 @@ import
org.apache.shardingsphere.infra.util.exception.ShardingSpherePrecondition
import
org.apache.shardingsphere.infra.util.exception.external.sql.type.wrapper.SQLWrapperException;
import java.sql.SQLException;
-import java.util.Arrays;
-import java.util.HashSet;
import java.util.LinkedHashMap;
+import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -58,16 +58,12 @@ public final class MigrationDataConsistencyChecker
implements PipelineDataConsis
private final JobRateLimitAlgorithm readRateLimitAlgorithm;
- private final TableNameSchemaNameMapping tableNameSchemaNameMapping;
-
private final ConsistencyCheckJobItemProgressContext progressContext;
public MigrationDataConsistencyChecker(final MigrationJobConfiguration
jobConfig, final InventoryIncrementalProcessContext processContext,
final
ConsistencyCheckJobItemProgressContext progressContext) {
this.jobConfig = jobConfig;
readRateLimitAlgorithm = null == processContext ? null :
processContext.getReadRateLimitAlgorithm();
- tableNameSchemaNameMapping = new TableNameSchemaNameMapping(
-
TableNameSchemaNameMapping.convert(jobConfig.getSourceSchemaName(), new
HashSet<>(Arrays.asList(jobConfig.getSourceTableName(),
jobConfig.getTargetTableName()))));
this.progressContext = progressContext;
}
@@ -75,8 +71,8 @@ public final class MigrationDataConsistencyChecker implements
PipelineDataConsis
public Map<String, DataConsistencyCheckResult> check(final
DataConsistencyCalculateAlgorithm calculateAlgorithm) {
verifyPipelineDatabaseType(calculateAlgorithm, jobConfig.getSource());
verifyPipelineDatabaseType(calculateAlgorithm, jobConfig.getTarget());
- SchemaTableName sourceTable = new SchemaTableName(new
SchemaName(tableNameSchemaNameMapping.getSchemaName(jobConfig.getSourceTableName())),
new TableName(jobConfig.getSourceTableName()));
- SchemaTableName targetTable = new SchemaTableName(new
SchemaName(tableNameSchemaNameMapping.getSchemaName(jobConfig.getTargetTableName())),
new TableName(jobConfig.getTargetTableName()));
+ SchemaTableName sourceTable = new SchemaTableName(new
SchemaName(jobConfig.getSourceSchemaName()), new
TableName(jobConfig.getSourceTableName()));
+ SchemaTableName targetTable = new SchemaTableName(new
SchemaName(jobConfig.getSourceSchemaName()), new
TableName(jobConfig.getTargetTableName()));
progressContext.getTableNames().add(jobConfig.getSourceTableName());
Map<String, DataConsistencyCheckResult> result = new LinkedHashMap<>();
try (
@@ -84,8 +80,11 @@ public final class MigrationDataConsistencyChecker
implements PipelineDataConsis
PipelineDataSourceWrapper targetDataSource =
PipelineDataSourceFactory.newInstance(jobConfig.getTarget())) {
progressContext.setRecordsCount(getRecordsCount());
PipelineTableMetaDataLoader metaDataLoader = new
StandardPipelineTableMetaDataLoader(sourceDataSource);
+ List<PipelineColumnMetaData> uniqueKeyColumns =
PipelineTableMetaDataUtil.getUniqueKeyColumns(
+ sourceTable.getSchemaName().getOriginal(),
sourceTable.getTableName().getOriginal(), metaDataLoader);
+ PipelineColumnMetaData uniqueKey = uniqueKeyColumns.isEmpty() ?
null : uniqueKeyColumns.get(0);
SingleTableInventoryDataConsistencyChecker
singleTableInventoryChecker = new SingleTableInventoryDataConsistencyChecker(
- jobConfig.getJobId(), sourceDataSource, targetDataSource,
sourceTable, targetTable, jobConfig.getUniqueKeyColumn(), metaDataLoader,
readRateLimitAlgorithm, progressContext);
+ jobConfig.getJobId(), sourceDataSource, targetDataSource,
sourceTable, targetTable, uniqueKey, metaDataLoader, readRateLimitAlgorithm,
progressContext);
result.put(sourceTable.getTableName().getOriginal(),
singleTableInventoryChecker.check(calculateAlgorithm));
} catch (final SQLException ex) {
throw new SQLWrapperException(ex);
diff --git
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/prepare/MigrationJobPreparer.java
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/prepare/MigrationJobPreparer.java
index ba4d7c9821c..1cb2c4114e3 100644
---
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/prepare/MigrationJobPreparer.java
+++
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/prepare/MigrationJobPreparer.java
@@ -148,10 +148,6 @@ public final class MigrationJobPreparer {
private void initInventoryTasks(final MigrationJobItemContext
jobItemContext) {
InventoryDumperConfiguration inventoryDumperConfig = new
InventoryDumperConfiguration(jobItemContext.getTaskConfig().getDumperConfig());
-
Optional.ofNullable(jobItemContext.getJobConfig().getUniqueKeyColumn()).ifPresent(uniqueKeyColumn
-> {
- inventoryDumperConfig.setUniqueKey(uniqueKeyColumn.getName());
-
inventoryDumperConfig.setUniqueKeyDataType(uniqueKeyColumn.getDataType());
- });
InventoryTaskSplitter inventoryTaskSplitter = new
InventoryTaskSplitter(jobItemContext.getSourceDataSource(),
inventoryDumperConfig, jobItemContext.getTaskConfig().getImporterConfig());
jobItemContext.getInventoryTasks().addAll(inventoryTaskSplitter.splitInventoryData(jobItemContext));
}
diff --git
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/IndexesMigrationE2EIT.java
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/IndexesMigrationE2EIT.java
index 56145d8ec51..eef05692249 100644
---
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/IndexesMigrationE2EIT.java
+++
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/IndexesMigrationE2EIT.java
@@ -90,6 +90,32 @@ public final class IndexesMigrationE2EIT extends
AbstractMigrationE2EIT {
assertMigrationSuccess(sql, consistencyCheckAlgorithmType);
}
+ @Test
+ public void assertMultiPrimaryKeyMigrationSuccess() throws SQLException,
InterruptedException {
+ String sql;
+ String consistencyCheckAlgorithmType;
+ if (getDatabaseType() instanceof MySQLDatabaseType) {
+ sql = "CREATE TABLE `%s` (`order_id` VARCHAR(64) NOT NULL,
`user_id` INT NOT NULL, `status` varchar(255), PRIMARY KEY
(`order_id`,`user_id`)) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4";
+ consistencyCheckAlgorithmType = "CRC32_MATCH";
+ } else {
+ return;
+ }
+ assertMigrationSuccess(sql, consistencyCheckAlgorithmType);
+ }
+
+ @Test
+ public void assertMultiUniqueKeyMigrationSuccess() throws SQLException,
InterruptedException {
+ String sql;
+ String consistencyCheckAlgorithmType;
+ if (getDatabaseType() instanceof MySQLDatabaseType) {
+ sql = "CREATE TABLE `%s` (`order_id` VARCHAR(64) NOT NULL,
`user_id` INT NOT NULL, `status` varchar(255), UNIQUE KEY
(`order_id`,`user_id`)) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4";
+ consistencyCheckAlgorithmType = "DATA_MATCH";
+ } else {
+ return;
+ }
+ assertMigrationSuccess(sql, consistencyCheckAlgorithmType);
+ }
+
@Test
public void assertSpecialTypeSingleColumnUniqueKeyMigrationSuccess()
throws SQLException, InterruptedException {
String sql;
diff --git
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/TextPrimaryKeyMigrationE2EIT.java
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/TextPrimaryKeyMigrationE2EIT.java
index 0cbac5e7961..559053f34a8 100644
---
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/TextPrimaryKeyMigrationE2EIT.java
+++
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/primarykey/TextPrimaryKeyMigrationE2EIT.java
@@ -70,7 +70,6 @@ public class TextPrimaryKeyMigrationE2EIT extends
AbstractMigrationE2EIT {
}
for (String version :
PipelineBaseE2EIT.ENV.listStorageContainerImages(new MySQLDatabaseType())) {
result.add(new PipelineTestParameter(new MySQLDatabaseType(),
version, "env/scenario/primary_key/text_primary_key/mysql.xml"));
- result.add(new PipelineTestParameter(new MySQLDatabaseType(),
version, "env/scenario/primary_key/unique_key/mysql.xml"));
}
for (String version :
PipelineBaseE2EIT.ENV.listStorageContainerImages(new PostgreSQLDatabaseType()))
{
result.add(new PipelineTestParameter(new PostgreSQLDatabaseType(),
version, "env/scenario/primary_key/text_primary_key/postgresql.xml"));
diff --git
a/test/e2e/pipeline/src/test/resources/env/scenario/primary_key/unique_key/mysql.xml
b/test/e2e/pipeline/src/test/resources/env/scenario/primary_key/unique_key/mysql.xml
deleted file mode 100644
index baeb0d95b75..00000000000
---
a/test/e2e/pipeline/src/test/resources/env/scenario/primary_key/unique_key/mysql.xml
+++ /dev/null
@@ -1,28 +0,0 @@
-<!--
- ~ Licensed to the Apache Software Foundation (ASF) under one or more
- ~ contributor license agreements. See the NOTICE file distributed with
- ~ this work for additional information regarding copyright ownership.
- ~ The ASF licenses this file to You under the Apache License, Version 2.0
- ~ (the "License"); you may not use this file except in compliance with
- ~ the License. You may obtain a copy of the License at
- ~
- ~ http://www.apache.org/licenses/LICENSE-2.0
- ~
- ~ Unless required by applicable law or agreed to in writing, software
- ~ distributed under the License is distributed on an "AS IS" BASIS,
- ~ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- ~ See the License for the specific language governing permissions and
- ~ limitations under the License.
- -->
-<command>
- <create-table-order>
- CREATE TABLE `%s` (
- `order_id` varchar(255) NOT NULL,
- `user_id` INT NOT NULL,
- `status` varchar(255) NULL,
- `t_unsigned_int` int UNSIGNED NULL,
- CONSTRAINT unique_id UNIQUE (order_id),
- INDEX ( `user_id` )
- ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
- </create-table-order>
-</command>
diff --git
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/impl/GovernanceRepositoryAPIImplTest.java
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/impl/GovernanceRepositoryAPIImplTest.java
index a836881ca50..f544f69e9ee 100644
---
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/impl/GovernanceRepositoryAPIImplTest.java
+++
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/impl/GovernanceRepositoryAPIImplTest.java
@@ -44,7 +44,7 @@ import
org.apache.shardingsphere.test.it.data.pipeline.core.util.PipelineContext
import org.junit.BeforeClass;
import org.junit.Test;
-import java.sql.Types;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -148,8 +148,7 @@ public final class GovernanceRepositoryAPIImplTest {
dumperConfig.setPosition(new PlaceholderPosition());
dumperConfig.setActualTableName("t_order");
dumperConfig.setLogicTableName("t_order");
- dumperConfig.setUniqueKey("order_id");
- dumperConfig.setUniqueKeyDataType(Types.INTEGER);
+
dumperConfig.setUniqueKeyColumns(Collections.singletonList(PipelineContextUtil.mockOrderIdColumnMetaData()));
dumperConfig.setShardingItem(0);
PipelineDataSourceWrapper dataSource =
mock(PipelineDataSourceWrapper.class);
PipelineTableMetaDataLoader metaDataLoader = new
StandardPipelineTableMetaDataLoader(dataSource);
diff --git
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/prepare/InventoryTaskSplitterTest.java
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/prepare/InventoryTaskSplitterTest.java
index 58f4796b3f1..e35557790df 100644
---
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/prepare/InventoryTaskSplitterTest.java
+++
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/prepare/InventoryTaskSplitterTest.java
@@ -28,7 +28,6 @@ import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.PipelineTabl
import
org.apache.shardingsphere.data.pipeline.core.metadata.loader.StandardPipelineTableMetaDataLoader;
import
org.apache.shardingsphere.data.pipeline.core.prepare.InventoryTaskSplitter;
import org.apache.shardingsphere.data.pipeline.core.task.InventoryTask;
-import
org.apache.shardingsphere.data.pipeline.scenario.migration.config.MigrationTaskConfiguration;
import
org.apache.shardingsphere.data.pipeline.scenario.migration.context.MigrationJobItemContext;
import
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
import
org.apache.shardingsphere.test.it.data.pipeline.core.util.PipelineContextUtil;
@@ -36,26 +35,24 @@ import org.junit.After;
import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.Test;
-import org.mockito.internal.configuration.plugins.Plugins;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Statement;
import java.sql.Types;
+import java.util.Collections;
import java.util.List;
-import java.util.Optional;
import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
-import static org.junit.Assert.assertFalse;
-import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
public final class InventoryTaskSplitterTest {
private MigrationJobItemContext jobItemContext;
- private MigrationTaskConfiguration taskConfig;
+ private InventoryDumperConfiguration dumperConfig;
private PipelineDataSourceManager dataSourceManager;
@@ -69,18 +66,16 @@ public final class InventoryTaskSplitterTest {
@Before
public void setUp() throws ReflectiveOperationException {
initJobItemContext();
- InventoryDumperConfiguration dumperConfig = new
InventoryDumperConfiguration(jobItemContext.getTaskConfig().getDumperConfig());
- dumperConfig.setUniqueKeyDataType(Types.INTEGER);
- dumperConfig.setUniqueKey("order_id");
+ dumperConfig = new
InventoryDumperConfiguration(jobItemContext.getTaskConfig().getDumperConfig());
+ PipelineColumnMetaData columnMetaData = new PipelineColumnMetaData(1,
"order_id", Types.INTEGER, "int", false, true, true);
+
dumperConfig.setUniqueKeyColumns(Collections.singletonList(columnMetaData));
inventoryTaskSplitter = new
InventoryTaskSplitter(jobItemContext.getSourceDataSource(), dumperConfig,
jobItemContext.getTaskConfig().getImporterConfig());
}
- private void initJobItemContext() throws ReflectiveOperationException {
+ private void initJobItemContext() {
MigrationJobConfiguration jobConfig =
JobConfigurationBuilder.createJobConfiguration();
-
Plugins.getMemberAccessor().set(MigrationJobConfiguration.class.getDeclaredField("uniqueKeyColumn"),
jobConfig, new PipelineColumnMetaData(1, "order_id", 4, "", false, true,
true));
jobItemContext =
PipelineContextUtil.mockMigrationJobItemContext(jobConfig);
dataSourceManager = (PipelineDataSourceManager)
jobItemContext.getImporterConnector().getConnector();
- taskConfig = jobItemContext.getTaskConfig();
}
@After
@@ -90,7 +85,7 @@ public final class InventoryTaskSplitterTest {
@Test
public void assertSplitInventoryDataWithEmptyTable() throws SQLException {
- initEmptyTablePrimaryEnvironment(taskConfig.getDumperConfig());
+ initEmptyTablePrimaryEnvironment(dumperConfig);
List<InventoryTask> actual =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
assertThat(actual.size(), is(1));
InventoryTask task = actual.get(0);
@@ -100,7 +95,7 @@ public final class InventoryTaskSplitterTest {
@Test
public void assertSplitInventoryDataWithIntPrimary() throws SQLException {
- initIntPrimaryEnvironment(taskConfig.getDumperConfig());
+ initIntPrimaryEnvironment(dumperConfig);
List<InventoryTask> actual =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
assertThat(actual.size(), is(10));
InventoryTask task = actual.get(9);
@@ -110,36 +105,34 @@ public final class InventoryTaskSplitterTest {
@Test
public void assertSplitInventoryDataWithCharPrimary() throws SQLException {
- initCharPrimaryEnvironment(taskConfig.getDumperConfig());
+ initCharPrimaryEnvironment(dumperConfig);
inventoryTaskSplitter.splitInventoryData(jobItemContext);
}
@Test
public void assertSplitInventoryDataWithoutPrimaryButWithUniqueIndex()
throws SQLException {
-
initUniqueIndexOnNotNullColumnEnvironment(taskConfig.getDumperConfig());
+ initUniqueIndexOnNotNullColumnEnvironment(dumperConfig);
List<InventoryTask> actual =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
assertThat(actual.size(), is(1));
}
@Test
- public void assertSplitInventoryDataWithMultipleColumnsKey() throws
SQLException, ReflectiveOperationException {
- initUnionPrimaryEnvironment(taskConfig.getDumperConfig());
- InventoryDumperConfiguration dumperConfig =
(InventoryDumperConfiguration) Plugins.getMemberAccessor()
-
.get(InventoryTaskSplitter.class.getDeclaredField("dumperConfig"),
inventoryTaskSplitter);
- assertNotNull(dumperConfig);
- dumperConfig.setUniqueKey("order_id,user_id");
- dumperConfig.setUniqueKeyDataType(Integer.MIN_VALUE);
- List<InventoryTask> actual =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
- assertThat(actual.size(), is(1));
+ public void assertSplitInventoryDataWithMultipleColumnsKey() throws
SQLException {
+ initUnionPrimaryEnvironment(dumperConfig);
+ try (PipelineDataSourceWrapper dataSource =
dataSourceManager.getDataSource(dumperConfig.getDataSourceConfig())) {
+ List<PipelineColumnMetaData> uniqueKeyColumns =
PipelineTableMetaDataUtil.getUniqueKeyColumns(null, "t_order", new
StandardPipelineTableMetaDataLoader(dataSource));
+ dumperConfig.setUniqueKeyColumns(uniqueKeyColumns);
+ List<InventoryTask> actual =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
+ assertThat(actual.size(), is(1));
+ }
}
@Test
- public void assertSplitInventoryDataWithoutPrimaryAndUniqueIndex() throws
SQLException, ReflectiveOperationException {
- initNoPrimaryEnvironment(taskConfig.getDumperConfig());
- try (PipelineDataSourceWrapper dataSource =
dataSourceManager.getDataSource(taskConfig.getDumperConfig().getDataSourceConfig()))
{
- Optional<PipelineColumnMetaData> uniqueKeyColumn =
PipelineTableMetaDataUtil.getUniqueKeyColumn(null, "t_order", new
StandardPipelineTableMetaDataLoader(dataSource));
- assertFalse(uniqueKeyColumn.isPresent());
-
Plugins.getMemberAccessor().set(MigrationJobConfiguration.class.getDeclaredField("uniqueKeyColumn"),
jobItemContext.getJobConfig(), null);
+ public void assertSplitInventoryDataWithoutPrimaryAndUniqueIndex() throws
SQLException {
+ initNoPrimaryEnvironment(dumperConfig);
+ try (PipelineDataSourceWrapper dataSource =
dataSourceManager.getDataSource(dumperConfig.getDataSourceConfig())) {
+ List<PipelineColumnMetaData> uniqueKeyColumns =
PipelineTableMetaDataUtil.getUniqueKeyColumns(null, "t_order", new
StandardPipelineTableMetaDataLoader(dataSource));
+ assertTrue(uniqueKeyColumns.isEmpty());
List<InventoryTask> inventoryTasks =
inventoryTaskSplitter.splitInventoryData(jobItemContext);
assertThat(inventoryTasks.size(), is(1));
}
diff --git
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/InventoryTaskTest.java
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/InventoryTaskTest.java
index e3ffd00534b..29db981d435 100644
---
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/InventoryTaskTest.java
+++
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/InventoryTaskTest.java
@@ -39,7 +39,7 @@ import org.junit.Test;
import java.sql.Connection;
import java.sql.SQLException;
import java.sql.Statement;
-import java.sql.Types;
+import java.util.Collections;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
@@ -115,8 +115,7 @@ public final class InventoryTaskTest {
InventoryDumperConfiguration result = new
InventoryDumperConfiguration(taskConfig.getDumperConfig());
result.setLogicTableName(logicTableName);
result.setActualTableName(actualTableName);
- result.setUniqueKey("order_id");
- result.setUniqueKeyDataType(Types.INTEGER);
+
result.setUniqueKeyColumns(Collections.singletonList(PipelineContextUtil.mockOrderIdColumnMetaData()));
result.setPosition(null == taskConfig.getDumperConfig().getPosition()
? new IntegerPrimaryKeyPosition(0, 1000) :
taskConfig.getDumperConfig().getPosition());
return result;
}
diff --git
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/JobConfigurationBuilder.java
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/JobConfigurationBuilder.java
index 7796e0fc0f0..73f2d2eb0d7 100644
---
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/JobConfigurationBuilder.java
+++
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/JobConfigurationBuilder.java
@@ -30,7 +30,6 @@ import
org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJobId
import
org.apache.shardingsphere.data.pipeline.scenario.migration.api.impl.MigrationJobAPI;
import
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigurationSwapper;
-import
org.apache.shardingsphere.data.pipeline.yaml.metadata.YamlPipelineColumnMetaData;
import org.apache.shardingsphere.infra.util.spi.type.typed.TypedSPILoader;
/**
@@ -54,23 +53,10 @@ public final class JobConfigurationBuilder {
result.setSource(createYamlPipelineDataSourceConfiguration(new
StandardPipelineDataSourceConfiguration(ConfigurationFileUtil.readFile("migration_standard_jdbc_source.yaml"))));
result.setTarget(createYamlPipelineDataSourceConfiguration(new
ShardingSpherePipelineDataSourceConfiguration(
ConfigurationFileUtil.readFile("migration_sharding_sphere_jdbc_target.yaml"))));
- result.setUniqueKeyColumn(createYamlPipelineColumnMetaData());
TypedSPILoader.getService(PipelineJobAPI.class,
"MIGRATION").extendYamlJobConfiguration(result);
return new YamlMigrationJobConfigurationSwapper().swapToObject(result);
}
- private static YamlPipelineColumnMetaData
createYamlPipelineColumnMetaData() {
- YamlPipelineColumnMetaData result = new YamlPipelineColumnMetaData();
- result.setOrdinalPosition(1);
- result.setName("order_id");
- result.setDataType(4);
- result.setDataTypeName("");
- result.setNullable(false);
- result.setPrimaryKey(true);
- result.setNullable(true);
- return result;
- }
-
private static String generateJobId(final YamlMigrationJobConfiguration
yamlJobConfig) {
String sourceTableName = RandomStringUtils.randomAlphabetic(32);
MigrationJobId migrationJobId = new
MigrationJobId(yamlJobConfig.getSourceResourceName(),
yamlJobConfig.getSourceSchemaName(), sourceTableName,
diff --git
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/PipelineContextUtil.java
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/PipelineContextUtil.java
index 8d8404bf030..e168cda57c9 100644
---
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/PipelineContextUtil.java
+++
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/util/PipelineContextUtil.java
@@ -21,6 +21,7 @@ import lombok.SneakyThrows;
import
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
import
org.apache.shardingsphere.data.pipeline.api.config.process.PipelineProcessConfiguration;
import
org.apache.shardingsphere.data.pipeline.api.datasource.config.impl.ShardingSpherePipelineDataSourceConfiguration;
+import
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumnMetaData;
import
org.apache.shardingsphere.data.pipeline.core.config.process.PipelineProcessConfigurationUtil;
import org.apache.shardingsphere.data.pipeline.core.context.PipelineContext;
import
org.apache.shardingsphere.data.pipeline.core.datasource.DefaultPipelineDataSourceManager;
@@ -101,6 +102,15 @@ public final class PipelineContextUtil {
return new MetaDataContexts(persistService, old.getMetaData());
}
+ /**
+ * Mock order_id column meta data.
+ *
+ * @return mocked column meta data
+ */
+ public static PipelineColumnMetaData mockOrderIdColumnMetaData() {
+ return new PipelineColumnMetaData(1, "order_id", Types.INTEGER, "int",
false, true, true);
+ }
+
/**
* Get execute engine.
*