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.
      *

Reply via email to