This is an automated email from the ASF dual-hosted git repository.

panjuan 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 17221f1335f Support migrating several tables at a time (#24342)
17221f1335f is described below

commit 17221f1335f5557e524e994a2fe9a84ef560a456
Author: Hongsheng Zhong <[email protected]>
AuthorDate: Fri Feb 24 20:37:40 2023 +0800

    Support migrating several tables at a time (#24342)
---
 .../user-manual/error-code/sql-error-code.cn.md    |   1 +
 .../user-manual/error-code/sql-error-code.en.md    |   3 +-
 .../shardingsphere/infra/datanode/DataNode.java    |   9 -
 .../infra/datanode/DataNodeTest.java               |  16 +-
 .../pipeline/api/config/ImporterConfiguration.java |   2 +-
 .../api/config/TableNameSchemaNameMapping.java     |  41 +--
 .../api/config/job/MigrationJobConfiguration.java  |  32 +--
 .../data/pipeline/api/datanode/DataNodeUtil.java   |  62 +++++
 .../pipeline/api/datanode/JobDataNodeEntry.java    |  18 +-
 .../pipeline/api/datanode/JobDataNodeLine.java     |   8 +-
 .../spi/check/datasource/DataSourceChecker.java    |   1 +
 .../api/config/TableNameSchemaNameMappingTest.java |  32 ++-
 .../pipeline/api/datanode/DataNodeUtilTest.java}   |  22 +-
 .../data/pipeline/cdc/api/impl/CDCJobAPI.java      |   4 +-
 .../param/PipelineInvalidParameterException.java}  |  25 +-
 .../handler/update/MigrateTableUpdater.java        |   4 +-
 .../core/MigrationDistSQLStatementVisitor.java     |   7 +-
 .../distsql/statement/MigrateTableStatement.java   |  11 +-
 .../distsql/statement/pojo/SourceTargetEntry.java} |  33 ++-
 kernel/data-pipeline/scenario/migration/pom.xml    |   5 +
 .../scenario/migration/MigrationJobId.java         |  19 +-
 .../migration/api/impl/MigrationJobAPI.java        | 276 ++++++++++++---------
 .../MigrationDataConsistencyChecker.java           |  63 +++--
 .../migration/prepare/MigrationJobPreparer.java    |  12 +-
 .../yaml/job/YamlMigrationJobConfiguration.java    |  33 +--
 .../job/YamlMigrationJobConfigurationSwapper.java  |  29 ++-
 .../pipeline/core/job/PipelineJobIdUtilsTest.java  |   4 +-
 .../general/MySQLMigrationGeneralE2EIT.java        |   4 +-
 .../general/PostgreSQLMigrationGeneralE2EIT.java   |   4 +-
 .../src/test/resources/env/mysql/01-initdb.sql     |  12 +-
 .../test/resources/env/postgresql/01-initdb.sql    |  10 +-
 .../update/MigrateTableStatementAssert.java        |  13 +-
 .../api/datanode/JobDataNodeEntryTest.java         |  17 ++
 .../DefaultPipelineDataSourceManagerTest.java      |   8 +-
 .../pipeline/core/task/IncrementalTaskTest.java    |   2 +-
 .../core/util/JobConfigurationBuilder.java         |  37 ++-
 .../migration/api/impl/MigrationJobAPITest.java    |  35 ++-
 .../MigrationDataConsistencyCheckerTest.java       |  13 +-
 .../YamlMigrationJobConfigurationSwapperTest.java  |  16 +-
 39 files changed, 553 insertions(+), 390 deletions(-)

diff --git a/docs/document/content/user-manual/error-code/sql-error-code.cn.md 
b/docs/document/content/user-manual/error-code/sql-error-code.cn.md
index 0820b5e9950..65c504568aa 100644
--- a/docs/document/content/user-manual/error-code/sql-error-code.cn.md
+++ b/docs/document/content/user-manual/error-code/sql-error-code.cn.md
@@ -107,6 +107,7 @@ SQL 错误码以标准的 SQL State,Vendor Code 和详细错误信息提供,
 | 42S02     | 18002       | There is no rule in database \`%s\`. |
 | 44000     | 18003       | Mode configuration does not exist. |
 | 44000     | 18004       | Target database name is null. You could define it 
in DistSQL or select a database. |
+| 22023     | 18005       | There is invalid parameter value: `%s`. |
 | HY000     | 18020       | Failed to get DDL for table \`%s\`. |
 | 42S01     | 18030       | Duplicate storage unit names \`%s\`. |
 | 42S02     | 18031       | Storage units names \`%s\` do not exist. |
diff --git a/docs/document/content/user-manual/error-code/sql-error-code.en.md 
b/docs/document/content/user-manual/error-code/sql-error-code.en.md
index 16689a17871..cc2043bed40 100644
--- a/docs/document/content/user-manual/error-code/sql-error-code.en.md
+++ b/docs/document/content/user-manual/error-code/sql-error-code.en.md
@@ -103,10 +103,11 @@ SQL error codes provide by standard `SQL State`, `Vendor 
Code` and `Reason`, whi
 ### Migration
 
 | SQL State | Vendor Code | Reason |
-| --------- | ----------- | ------ |
+| --------- |-------------| ------ |
 | 42S02     | 18002       | There is no rule in database \`%s\`. |
 | 44000     | 18003       | Mode configuration does not exist. |
 | 44000     | 18004       | Target database name is null. You could define it 
in DistSQL or select a database. |
+| 22023     | 18005       | There is invalid parameter value: `%s`. |
 | HY000     | 18020       | Failed to get DDL for table \`%s\`. |
 | 42S01     | 18030       | Duplicate storage unit names \`%s\`. |
 | 42S02     | 18031       | Storage units names \`%s\` do not exist. |
diff --git 
a/infra/common/src/main/java/org/apache/shardingsphere/infra/datanode/DataNode.java
 
b/infra/common/src/main/java/org/apache/shardingsphere/infra/datanode/DataNode.java
index 56706d797da..0bb17be5203 100644
--- 
a/infra/common/src/main/java/org/apache/shardingsphere/infra/datanode/DataNode.java
+++ 
b/infra/common/src/main/java/org/apache/shardingsphere/infra/datanode/DataNode.java
@@ -96,15 +96,6 @@ public final class DataNode {
         return dataSourceName + DELIMITER + tableName;
     }
     
-    /**
-     * Get formatted text length.
-     *
-     * @return formatted text length
-     */
-    public int getFormattedTextLength() {
-        return dataSourceName.length() + DELIMITER.length() + 
tableName.length();
-    }
-    
     /**
      * Is Actual data nodes three tier structure.
      *
diff --git 
a/infra/common/src/test/java/org/apache/shardingsphere/infra/datanode/DataNodeTest.java
 
b/infra/common/src/test/java/org/apache/shardingsphere/infra/datanode/DataNodeTest.java
index 9b4a63a31b2..dc3108ec74f 100644
--- 
a/infra/common/src/test/java/org/apache/shardingsphere/infra/datanode/DataNodeTest.java
+++ 
b/infra/common/src/test/java/org/apache/shardingsphere/infra/datanode/DataNodeTest.java
@@ -22,8 +22,8 @@ import org.junit.Test;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.CoreMatchers.not;
-import static org.junit.Assert.assertFalse;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.Assert.assertFalse;
 
 public final class DataNodeTest {
     
@@ -85,13 +85,6 @@ public final class DataNodeTest {
         assertThat(dataNode.format(), is(expected));
     }
     
-    @Test
-    public void assertFormattedTextLength() {
-        String text = "ds_0.tbl_0";
-        DataNode dataNode = new DataNode(text);
-        assertThat(dataNode.getFormattedTextLength(), is(text.length()));
-    }
-    
     @Test
     public void assertNewValidDataNodeIncludeInstance() {
         DataNode dataNode = new DataNode("ds_0.db_0.tbl_0");
@@ -115,11 +108,4 @@ public final class DataNodeTest {
         DataNode dataNode = new DataNode(expected);
         assertThat(dataNode.format(), is(expected));
     }
-    
-    @Test
-    public void assertFormattedTextLengthIncludeInstance() {
-        String text = "ds_0.db_0.tbl_0";
-        DataNode dataNode = new DataNode(text);
-        assertThat(dataNode.getFormattedTextLength(), is(text.length()));
-    }
 }
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ImporterConfiguration.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ImporterConfiguration.java
index 7af6e8d9c3a..7740caff919 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ImporterConfiguration.java
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/ImporterConfiguration.java
@@ -37,7 +37,7 @@ import java.util.stream.Collectors;
  */
 @RequiredArgsConstructor
 @Getter
-@ToString(exclude = "dataSourceConfig")
+@ToString(exclude = {"dataSourceConfig", "tableNameSchemaNameMapping"})
 public final class ImporterConfiguration {
     
     private final PipelineDataSourceConfiguration dataSourceConfig;
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMapping.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMapping.java
index 9b1de2ab1d8..160ba945b41 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMapping.java
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMapping.java
@@ -17,19 +17,16 @@
 
 package org.apache.shardingsphere.data.pipeline.api.config;
 
-import lombok.RequiredArgsConstructor;
 import lombok.ToString;
 import org.apache.shardingsphere.data.pipeline.api.metadata.LogicTableName;
 
-import java.util.Collection;
-import java.util.LinkedHashMap;
-import java.util.List;
+import java.util.Collections;
+import java.util.HashMap;
 import java.util.Map;
 
 /**
  * Table name and schema name mapping.
  */
-@RequiredArgsConstructor
 @ToString
 public final class TableNameSchemaNameMapping {
     
@@ -38,32 +35,20 @@ public final class TableNameSchemaNameMapping {
     /**
      * Convert table name and schema name mapping from schemas.
      *
-     * @param schemaTablesMap schema name and table names map
-     * @return logic table name and schema name map
+     * @param tableSchemaMap table name and schema name map
      */
-    public static Map<LogicTableName, String> convert(final Map<String, 
List<String>> schemaTablesMap) {
-        Map<LogicTableName, String> result = new LinkedHashMap<>();
-        schemaTablesMap.forEach((schemaName, tableNames) -> {
-            for (String each : tableNames) {
-                result.put(new LogicTableName(each), schemaName);
+    public TableNameSchemaNameMapping(final Map<String, String> 
tableSchemaMap) {
+        if (null == tableSchemaMap) {
+            mapping = Collections.emptyMap();
+            return;
+        }
+        Map<LogicTableName, String> mapping = new HashMap<>();
+        tableSchemaMap.forEach((tableName, schemaName) -> {
+            if (null != schemaName) {
+                mapping.put(new LogicTableName(tableName), schemaName);
             }
         });
-        return result;
-    }
-    
-    /**
-     * Convert table name and schema name mapping.
-     *
-     * @param schemaName schema name
-     * @param tables tables
-     * @return logic table name and schema name map
-     */
-    public static Map<LogicTableName, String> convert(final String schemaName, 
final Collection<String> tables) {
-        Map<LogicTableName, String> result = new LinkedHashMap<>();
-        for (String each : tables) {
-            result.put(new LogicTableName(each), schemaName);
-        }
-        return result;
+        this.mapping = mapping;
     }
     
     /**
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 ddbd3478fe5..88865848d81 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
@@ -20,50 +20,42 @@ package 
org.apache.shardingsphere.data.pipeline.api.config.job;
 import lombok.Getter;
 import lombok.RequiredArgsConstructor;
 import lombok.ToString;
+import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeLine;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
 
 import java.util.List;
+import java.util.Map;
 
 /**
  * Migration job configuration.
  */
 @RequiredArgsConstructor
 @Getter
-@ToString(exclude = {"source", "target"})
+@ToString(exclude = {"sources", "target"})
 public final class MigrationJobConfiguration implements 
PipelineJobConfiguration {
     
     private final String jobId;
     
-    private final String sourceResourceName;
-    
     private final String targetDatabaseName;
     
-    private final String sourceSchemaName;
-    
-    // TODO add targetSchemaName
-    
     private final String sourceDatabaseType;
     
     private final String targetDatabaseType;
     
-    private final String sourceTableName;
-    
-    private final String targetTableName;
-    
-    private final PipelineDataSourceConfiguration source;
+    private final Map<String, PipelineDataSourceConfiguration> sources;
     
     private final PipelineDataSourceConfiguration target;
     
+    private final List<String> targetTableNames;
+    
     /**
-     * Collection of each logic table's first data node.
-     * <p>
-     * If <pre>actualDataNodes: ds_${0..1}.t_order_${0..1}</pre> and 
<pre>actualDataNodes: ds_${0..1}.t_order_item_${0..1}</pre>,
-     * then value may be: {@code 
t_order:ds_0.t_order_0|t_order_item:ds_0.t_order_item_0}.
-     * </p>
+     * Map{logic table names, schema name}.
      */
-    private final String tablesFirstDataNodes;
+    private final Map<String, String> targetTableSchemaMap;
+    
+    private final JobDataNodeLine tablesFirstDataNodes;
     
-    private final List<String> jobShardingDataNodes;
+    private final List<JobDataNodeLine> jobShardingDataNodes;
     
     private final int concurrency;
     
@@ -75,6 +67,6 @@ public final class MigrationJobConfiguration implements 
PipelineJobConfiguration
      * @return job sharding count
      */
     public int getJobShardingCount() {
-        return 1;
+        return jobShardingDataNodes.size();
     }
 }
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtil.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtil.java
new file mode 100644
index 00000000000..ecaa6a6c96e
--- /dev/null
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtil.java
@@ -0,0 +1,62 @@
+/*
+ * 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.api.datanode;
+
+import com.google.common.base.Splitter;
+import lombok.AccessLevel;
+import lombok.NoArgsConstructor;
+import org.apache.shardingsphere.infra.datanode.DataNode;
+import 
org.apache.shardingsphere.infra.exception.InvalidDataNodesFormatException;
+
+import java.util.List;
+
+/**
+ * Data node util.
+ */
+@NoArgsConstructor(access = AccessLevel.PRIVATE)
+public final class DataNodeUtil {
+    
+    /**
+     * Format data node as string with schema.
+     *
+     * @param dataNode data node
+     * @return formatted data node
+     */
+    public static String formatWithSchema(final DataNode dataNode) {
+        return dataNode.getDataSourceName() + (null != 
dataNode.getSchemaName() ? "." + dataNode.getSchemaName() : "") + "." + 
dataNode.getTableName();
+    }
+    
+    /**
+     * Parse data node from text.
+     *
+     * @param text data node text
+     * @return data node
+     */
+    public static DataNode parseWithSchema(final String text) {
+        List<String> segments = Splitter.on(".").splitToList(text);
+        boolean hasSchema = 3 == segments.size();
+        if (!(2 == segments.size() || hasSchema)) {
+            throw new InvalidDataNodesFormatException(text);
+        }
+        DataNode result = new DataNode(segments.get(0), 
segments.get(segments.size() - 1));
+        if (hasSchema) {
+            result.setSchemaName(segments.get(1));
+        }
+        return result;
+    }
+}
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeEntry.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeEntry.java
index fc5682c300d..ce9ada6d023 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeEntry.java
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeEntry.java
@@ -22,7 +22,6 @@ import lombok.Getter;
 import lombok.RequiredArgsConstructor;
 import org.apache.shardingsphere.infra.datanode.DataNode;
 
-import java.util.Collection;
 import java.util.LinkedList;
 import java.util.List;
 
@@ -35,7 +34,7 @@ public final class JobDataNodeEntry {
     
     private final String logicTableName;
     
-    private final Collection<DataNode> dataNodes;
+    private final List<DataNode> dataNodes;
     
     /**
      * Unmarshal from text.
@@ -48,7 +47,7 @@ public final class JobDataNodeEntry {
         String logicTableName = segments.get(0);
         List<DataNode> dataNodes = new LinkedList<>();
         for (String each : 
Splitter.on(",").omitEmptyStrings().splitToList(segments.get(1))) {
-            dataNodes.add(new DataNode(each));
+            dataNodes.add(DataNodeUtil.parseWithSchema(each));
         }
         return new JobDataNodeEntry(logicTableName, dataNodes);
     }
@@ -59,24 +58,15 @@ public final class JobDataNodeEntry {
      * @return text, format: logicTableName:dataNode1,dataNode2, e.g. 
t_order:ds_0.t_order_0,ds_0.t_order_1
      */
     public String marshal() {
-        StringBuilder result = new 
StringBuilder(getMarshalledTextEstimatedLength());
+        StringBuilder result = new StringBuilder();
         result.append(logicTableName);
         result.append(":");
         for (DataNode each : dataNodes) {
-            result.append(each.format()).append(',');
+            result.append(DataNodeUtil.formatWithSchema(each)).append(',');
         }
         if (!dataNodes.isEmpty()) {
             result.setLength(result.length() - 1);
         }
         return result.toString();
     }
-    
-    /**
-     * Get marshalled text estimated length.
-     *
-     * @return marshalled text estimated length
-     */
-    public int getMarshalledTextEstimatedLength() {
-        return logicTableName.length() + 1 + 
dataNodes.stream().mapToInt(DataNode::getFormattedTextLength).sum() + 
dataNodes.size();
-    }
 }
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeLine.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeLine.java
index d43785fbace..f49e5a8cd64 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeLine.java
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/datanode/JobDataNodeLine.java
@@ -21,7 +21,7 @@ import com.google.common.base.Splitter;
 import lombok.Getter;
 import lombok.RequiredArgsConstructor;
 
-import java.util.Collection;
+import java.util.List;
 import java.util.stream.Collectors;
 
 /**
@@ -29,9 +29,11 @@ import java.util.stream.Collectors;
  */
 @RequiredArgsConstructor
 @Getter
+// TODO Move to pipeline-core
 public final class JobDataNodeLine {
     
-    private final Collection<JobDataNodeEntry> entries;
+    // Need sequential collection
+    private final List<JobDataNodeEntry> entries;
     
     /**
      * Marshal to text.
@@ -39,7 +41,7 @@ public final class JobDataNodeLine {
      * @return marshaled text, format: entry1|entry2, e.g. 
t_order:ds_0.t_order_0,ds_0.t_order_1|t_order_item:ds_0.t_order_item_0,ds_0.t_order_item_1
      */
     public String marshal() {
-        StringBuilder result = new 
StringBuilder(entries.stream().mapToInt(JobDataNodeEntry::getMarshalledTextEstimatedLength).sum()
 + entries.size());
+        StringBuilder result = new StringBuilder();
         for (JobDataNodeEntry each : entries) {
             result.append(each.marshal()).append('|');
         }
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/check/datasource/DataSourceChecker.java
 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/check/datasource/DataSourceChecker.java
index ff3a18ef455..7f2ba48c3c1 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/check/datasource/DataSourceChecker.java
+++ 
b/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/spi/check/datasource/DataSourceChecker.java
@@ -59,5 +59,6 @@ public interface DataSourceChecker extends TypedSPI {
      * @param logicTableNames logic table names
      */
     // TODO rename to common usage name
+    // TODO Merge schemaName and tableNames
     void checkTargetTable(Collection<? extends DataSource> dataSources, 
TableNameSchemaNameMapping tableNameSchemaNameMapping, Collection<String> 
logicTableNames);
 }
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
 
b/kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMappingTest.java
similarity index 51%
copy from 
test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
copy to 
kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMappingTest.java
index 94423f7bfa4..192ec922db3 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
+++ 
b/kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/config/TableNameSchemaNameMappingTest.java
@@ -15,23 +15,37 @@
  * limitations under the License.
  */
 
-package org.apache.shardingsphere.test.it.data.pipeline.api.datanode;
+package org.apache.shardingsphere.data.pipeline.api.config;
 
-import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeEntry;
-import org.apache.shardingsphere.infra.datanode.DataNode;
 import org.junit.Test;
 
-import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.Assert.assertNull;
 
-public final class JobDataNodeEntryTest {
+public final class TableNameSchemaNameMappingTest {
     
     @Test
-    public void assertMarshal() {
-        String actual = new JobDataNodeEntry("t_order", Arrays.asList(new 
DataNode("ds_0.t_order_0"), new DataNode("ds_0.t_order_1"))).marshal();
-        String expected = "t_order:ds_0.t_order_0,ds_0.t_order_1";
-        assertThat(actual, is(expected));
+    public void assertConstructFromNull() {
+        new TableNameSchemaNameMapping(null);
+    }
+    
+    @Test
+    public void assertConstructFromValueNullMap() {
+        Map<String, String> map = new HashMap<>();
+        map.put("t_order", null);
+        TableNameSchemaNameMapping mapping = new 
TableNameSchemaNameMapping(map);
+        assertNull(mapping.getSchemaName("t_order"));
+    }
+    
+    @Test
+    public void assertConstructFromMap() {
+        Map<String, String> map = new HashMap<>();
+        map.put("t_order", "public");
+        TableNameSchemaNameMapping mapping = new 
TableNameSchemaNameMapping(map);
+        assertThat(mapping.getSchemaName("t_order"), is("public"));
     }
 }
diff --git 
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
 
b/kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtilTest.java
similarity index 57%
copy from 
kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
copy to 
kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtilTest.java
index 05e0f6b21e6..711082bf09b 100644
--- 
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
+++ 
b/kernel/data-pipeline/api/src/test/java/org/apache/shardingsphere/data/pipeline/api/datanode/DataNodeUtilTest.java
@@ -15,20 +15,28 @@
  * limitations under the License.
  */
 
-package org.apache.shardingsphere.data.pipeline.yaml.job;
+package org.apache.shardingsphere.data.pipeline.api.datanode;
 
-import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
+import org.apache.shardingsphere.infra.datanode.DataNode;
 import org.junit.Test;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.MatcherAssert.assertThat;
 
-public final class YamlMigrationJobConfigurationSwapperTest {
+public final class DataNodeUtilTest {
     
     @Test
-    public void assertSwapToObject() {
-        YamlMigrationJobConfiguration yamlJobConfig = new 
YamlMigrationJobConfiguration();
-        MigrationJobConfiguration actual = new 
YamlMigrationJobConfigurationSwapper().swapToObject(yamlJobConfig);
-        assertThat(actual.getJobShardingCount(), is(1));
+    public void assertFormatWithSchema() {
+        DataNode dataNode = new DataNode("ds_0.tbl_0");
+        dataNode.setSchemaName("public");
+        assertThat(DataNodeUtil.formatWithSchema(dataNode), 
is("ds_0.public.tbl_0"));
+    }
+    
+    @Test
+    public void assertParseWithSchema() {
+        DataNode actual = DataNodeUtil.parseWithSchema("ds_0.public.tbl_0");
+        assertThat(actual.getDataSourceName(), is("ds_0"));
+        assertThat(actual.getSchemaName(), is("public"));
+        assertThat(actual.getTableName(), is("tbl_0"));
     }
 }
diff --git 
a/kernel/data-pipeline/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/impl/CDCJobAPI.java
 
b/kernel/data-pipeline/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/impl/CDCJobAPI.java
index 4a3a3495cd7..76e1c96e84a 100644
--- 
a/kernel/data-pipeline/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/impl/CDCJobAPI.java
+++ 
b/kernel/data-pipeline/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/impl/CDCJobAPI.java
@@ -226,11 +226,11 @@ public final class CDCJobAPI extends 
AbstractInventoryIncrementalJobAPIImpl {
     }
     
     private TableNameSchemaNameMapping getTableNameSchemaNameMapping(final 
Collection<String> tableNames) {
-        Map<LogicTableName, String> tableNameSchemaMap = new LinkedHashMap<>();
+        Map<String, String> tableNameSchemaMap = new LinkedHashMap<>();
         for (String each : tableNames) {
             String[] split = each.split("\\.");
             if (split.length > 1) {
-                tableNameSchemaMap.put(new LogicTableName(split[1]), split[0]);
+                tableNameSchemaMap.put(split[1], split[0]);
             }
         }
         return new TableNameSchemaNameMapping(tableNameSchemaMap);
diff --git 
a/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
 
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/exception/param/PipelineInvalidParameterException.java
similarity index 56%
copy from 
kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
copy to 
kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/exception/param/PipelineInvalidParameterException.java
index ee214f4b5db..3802597e3ba 100644
--- 
a/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
+++ 
b/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/exception/param/PipelineInvalidParameterException.java
@@ -15,26 +15,19 @@
  * limitations under the License.
  */
 
-package org.apache.shardingsphere.migration.distsql.statement;
+package org.apache.shardingsphere.data.pipeline.core.exception.param;
 
-import lombok.Getter;
-import lombok.RequiredArgsConstructor;
-import 
org.apache.shardingsphere.distsql.parser.statement.ral.scaling.UpdatableScalingRALStatement;
+import 
org.apache.shardingsphere.data.pipeline.core.exception.PipelineSQLException;
+import 
org.apache.shardingsphere.infra.util.exception.external.sql.sqlstate.XOpenSQLState;
 
 /**
- * Migrate table statement.
+ * Pipeline invalid parameter exception.
  */
-@RequiredArgsConstructor
-@Getter
-public final class MigrateTableStatement extends UpdatableScalingRALStatement {
+public final class PipelineInvalidParameterException extends 
PipelineSQLException {
     
-    private final String sourceResourceName;
+    private static final long serialVersionUID = -2162309404414015630L;
     
-    private final String sourceSchemaName;
-    
-    private final String sourceTableName;
-    
-    private final String targetDatabaseName;
-    
-    private final String targetTableName;
+    public PipelineInvalidParameterException(final String message) {
+        super(XOpenSQLState.INVALID_PARAMETER_VALUE, 5, String.format("There 
is invalid parameter value: %s.", message));
+    }
 }
diff --git 
a/kernel/data-pipeline/distsql/handler/src/main/java/org/apache/shardingsphere/migration/distsql/handler/update/MigrateTableUpdater.java
 
b/kernel/data-pipeline/distsql/handler/src/main/java/org/apache/shardingsphere/migration/distsql/handler/update/MigrateTableUpdater.java
index f1cda7a3f67..b0f84811ca0 100644
--- 
a/kernel/data-pipeline/distsql/handler/src/main/java/org/apache/shardingsphere/migration/distsql/handler/update/MigrateTableUpdater.java
+++ 
b/kernel/data-pipeline/distsql/handler/src/main/java/org/apache/shardingsphere/migration/distsql/handler/update/MigrateTableUpdater.java
@@ -18,7 +18,6 @@
 package org.apache.shardingsphere.migration.distsql.handler.update;
 
 import lombok.extern.slf4j.Slf4j;
-import 
org.apache.shardingsphere.data.pipeline.api.pojo.CreateMigrationJobParameter;
 import 
org.apache.shardingsphere.data.pipeline.core.exception.job.MissingRequiredTargetDatabaseException;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.api.impl.MigrationJobAPI;
 import org.apache.shardingsphere.distsql.handler.ral.update.RALUpdater;
@@ -37,8 +36,7 @@ public final class MigrateTableUpdater implements 
RALUpdater<MigrateTableStateme
     public void executeUpdate(final String databaseName, final 
MigrateTableStatement sqlStatement) {
         String targetDatabaseName = null == 
sqlStatement.getTargetDatabaseName() ? databaseName : 
sqlStatement.getTargetDatabaseName();
         ShardingSpherePreconditions.checkNotNull(targetDatabaseName, 
MissingRequiredTargetDatabaseException::new);
-        jobAPI.createJobAndStart(new CreateMigrationJobParameter(
-                sqlStatement.getSourceResourceName(), 
sqlStatement.getSourceSchemaName(), sqlStatement.getSourceTableName(), 
targetDatabaseName, sqlStatement.getTargetTableName()));
+        jobAPI.createJobAndStart(new 
MigrateTableStatement(sqlStatement.getSourceTargetEntries(), 
targetDatabaseName));
     }
     
     @Override
diff --git 
a/kernel/data-pipeline/distsql/parser/src/main/java/org/apache/shardingsphere/migration/distsql/parser/core/MigrationDistSQLStatementVisitor.java
 
b/kernel/data-pipeline/distsql/parser/src/main/java/org/apache/shardingsphere/migration/distsql/parser/core/MigrationDistSQLStatementVisitor.java
index 2e22d740296..4721476bde0 100644
--- 
a/kernel/data-pipeline/distsql/parser/src/main/java/org/apache/shardingsphere/migration/distsql/parser/core/MigrationDistSQLStatementVisitor.java
+++ 
b/kernel/data-pipeline/distsql/parser/src/main/java/org/apache/shardingsphere/migration/distsql/parser/core/MigrationDistSQLStatementVisitor.java
@@ -46,6 +46,7 @@ import 
org.apache.shardingsphere.distsql.parser.segment.AlgorithmSegment;
 import org.apache.shardingsphere.distsql.parser.segment.DataSourceSegment;
 import 
org.apache.shardingsphere.distsql.parser.segment.HostnameAndPortBasedDataSourceSegment;
 import 
org.apache.shardingsphere.distsql.parser.segment.URLBasedDataSourceSegment;
+import org.apache.shardingsphere.infra.datanode.DataNode;
 import 
org.apache.shardingsphere.migration.distsql.statement.CheckMigrationStatement;
 import 
org.apache.shardingsphere.migration.distsql.statement.CommitMigrationStatement;
 import 
org.apache.shardingsphere.migration.distsql.statement.DropMigrationCheckStatement;
@@ -62,11 +63,13 @@ import 
org.apache.shardingsphere.migration.distsql.statement.StartMigrationState
 import 
org.apache.shardingsphere.migration.distsql.statement.StopMigrationCheckStatement;
 import 
org.apache.shardingsphere.migration.distsql.statement.StopMigrationStatement;
 import 
org.apache.shardingsphere.migration.distsql.statement.UnregisterMigrationSourceStorageUnitStatement;
+import 
org.apache.shardingsphere.migration.distsql.statement.pojo.SourceTargetEntry;
 import org.apache.shardingsphere.sql.parser.api.visitor.ASTNode;
 import org.apache.shardingsphere.sql.parser.api.visitor.SQLVisitor;
 import 
org.apache.shardingsphere.sql.parser.sql.common.value.identifier.IdentifierValue;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.LinkedList;
 import java.util.List;
 import java.util.Properties;
@@ -86,7 +89,9 @@ public final class MigrationDistSQLStatementVisitor extends 
MigrationDistSQLStat
         String sourceTableName = source.get(source.size() - 1);
         String targetDatabaseName = target.size() > 1 ? target.get(0) : null;
         String targetTableName = target.get(target.size() - 1);
-        return new MigrateTableStatement(sourceResourceName, sourceSchemaName, 
sourceTableName, targetDatabaseName, targetTableName);
+        SourceTargetEntry sourceTargetEntry = new 
SourceTargetEntry(targetDatabaseName, new DataNode(sourceResourceName, 
sourceTableName), targetTableName);
+        sourceTargetEntry.getSource().setSchemaName(sourceSchemaName);
+        return new 
MigrateTableStatement(Collections.singletonList(sourceTargetEntry), 
targetDatabaseName);
     }
     
     @Override
diff --git 
a/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
 
b/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
index ee214f4b5db..83b27d94aea 100644
--- 
a/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
+++ 
b/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/MigrateTableStatement.java
@@ -20,6 +20,9 @@ package org.apache.shardingsphere.migration.distsql.statement;
 import lombok.Getter;
 import lombok.RequiredArgsConstructor;
 import 
org.apache.shardingsphere.distsql.parser.statement.ral.scaling.UpdatableScalingRALStatement;
+import 
org.apache.shardingsphere.migration.distsql.statement.pojo.SourceTargetEntry;
+
+import java.util.List;
 
 /**
  * Migrate table statement.
@@ -28,13 +31,7 @@ import 
org.apache.shardingsphere.distsql.parser.statement.ral.scaling.UpdatableS
 @Getter
 public final class MigrateTableStatement extends UpdatableScalingRALStatement {
     
-    private final String sourceResourceName;
-    
-    private final String sourceSchemaName;
-    
-    private final String sourceTableName;
+    private final List<SourceTargetEntry> sourceTargetEntries;
     
     private final String targetDatabaseName;
-    
-    private final String targetTableName;
 }
diff --git 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/pojo/CreateMigrationJobParameter.java
 
b/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/pojo/SourceTargetEntry.java
similarity index 57%
rename from 
kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/pojo/CreateMigrationJobParameter.java
rename to 
kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/pojo/SourceTargetEntry.java
index 9629da29985..72c89b355bb 100644
--- 
a/kernel/data-pipeline/api/src/main/java/org/apache/shardingsphere/data/pipeline/api/pojo/CreateMigrationJobParameter.java
+++ 
b/kernel/data-pipeline/distsql/statement/src/main/java/org/apache/shardingsphere/migration/distsql/statement/pojo/SourceTargetEntry.java
@@ -15,22 +15,41 @@
  * limitations under the License.
  */
 
-package org.apache.shardingsphere.data.pipeline.api.pojo;
+package org.apache.shardingsphere.migration.distsql.statement.pojo;
 
 import lombok.Getter;
 import lombok.RequiredArgsConstructor;
+import org.apache.shardingsphere.infra.datanode.DataNode;
 
+import java.util.Objects;
+
+/**
+ * Source target entry.
+ */
 @RequiredArgsConstructor
 @Getter
-public final class CreateMigrationJobParameter {
+public final class SourceTargetEntry {
     
-    private final String sourceResourceName;
+    private final String targetDatabaseName;
     
-    private final String sourceSchemaName;
+    private final DataNode source;
     
-    private final String sourceTableName;
+    private final String targetTableName;
     
-    private final String targetDatabaseName;
+    @Override
+    public boolean equals(final Object object) {
+        if (this == object) {
+            return true;
+        }
+        if (null == object || getClass() != object.getClass()) {
+            return false;
+        }
+        final SourceTargetEntry that = (SourceTargetEntry) object;
+        return source.equals(that.source) && 
targetTableName.equals(that.targetTableName);
+    }
     
-    private final String targetTableName;
+    @Override
+    public int hashCode() {
+        return Objects.hash(source, targetTableName);
+    }
 }
diff --git a/kernel/data-pipeline/scenario/migration/pom.xml 
b/kernel/data-pipeline/scenario/migration/pom.xml
index 6405d6653e4..062f3cdd040 100644
--- a/kernel/data-pipeline/scenario/migration/pom.xml
+++ b/kernel/data-pipeline/scenario/migration/pom.xml
@@ -28,6 +28,11 @@
     <name>${project.artifactId}</name>
     
     <dependencies>
+        <dependency>
+            <groupId>org.apache.shardingsphere</groupId>
+            
<artifactId>shardingsphere-data-pipeline-distsql-statement</artifactId>
+            <version>${project.version}</version>
+        </dependency>
         <dependency>
             <groupId>org.apache.shardingsphere</groupId>
             <artifactId>shardingsphere-data-pipeline-core</artifactId>
diff --git 
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/MigrationJobId.java
 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/MigrationJobId.java
index 15cc3991430..c8c1ab2f5b3 100644
--- 
a/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/MigrationJobId.java
+++ 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/migration/MigrationJobId.java
@@ -18,10 +18,11 @@
 package org.apache.shardingsphere.data.pipeline.scenario.migration;
 
 import lombok.Getter;
-import lombok.NonNull;
 import lombok.ToString;
 import org.apache.shardingsphere.data.pipeline.core.job.AbstractPipelineJobId;
 
+import java.util.List;
+
 /**
  * Migration job id.
  */
@@ -31,23 +32,13 @@ public final class MigrationJobId extends 
AbstractPipelineJobId {
     
     public static final String CURRENT_VERSION = "01";
     
-    private final String sourceResourceName;
-    
-    private final String sourceSchemaName;
-    
-    private final String sourceTableName;
+    private final List<String> jobShardingDataNodes;
     
     private final String targetDatabaseName;
     
-    private final String targetTableName;
-    
-    public MigrationJobId(@NonNull final String sourceResourceName, final 
String sourceSchemaName, @NonNull final String sourceTableName,
-                          @NonNull final String targetDatabaseName, @NonNull 
final String targetTableName) {
+    public MigrationJobId(final List<String> jobShardingDataNodes, final 
String targetDatabaseName) {
         super(new MigrationJobType(), CURRENT_VERSION);
-        this.sourceResourceName = sourceResourceName;
-        this.sourceSchemaName = sourceSchemaName;
-        this.sourceTableName = sourceTableName;
+        this.jobShardingDataNodes = jobShardingDataNodes;
         this.targetDatabaseName = targetDatabaseName;
-        this.targetTableName = targetTableName;
     }
 }
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 95b9c112e47..73a9eab3fbd 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
@@ -17,8 +17,6 @@
 
 package org.apache.shardingsphere.data.pipeline.scenario.migration.api.impl;
 
-import com.google.common.base.Joiner;
-import com.google.common.base.Strings;
 import com.google.gson.Gson;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.commons.codec.digest.DigestUtils;
@@ -32,11 +30,11 @@ import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfig
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.PipelineJobConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.yaml.YamlPipelineJobConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.config.process.PipelineProcessConfiguration;
+import org.apache.shardingsphere.data.pipeline.api.datanode.DataNodeUtil;
 import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeEntry;
 import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeLine;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
-import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfigurationFactory;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.impl.ShardingSpherePipelineDataSourceConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.impl.StandardPipelineDataSourceConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.yaml.YamlPipelineDataSourceConfiguration;
@@ -46,7 +44,6 @@ import 
org.apache.shardingsphere.data.pipeline.api.metadata.LogicTableName;
 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.pojo.CreateMigrationJobParameter;
 import org.apache.shardingsphere.data.pipeline.api.pojo.PipelineJobMetaData;
 import 
org.apache.shardingsphere.data.pipeline.api.pojo.TableBasedPipelineJobInfo;
 import org.apache.shardingsphere.data.pipeline.core.api.PipelineAPIFactory;
@@ -60,8 +57,10 @@ 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.exception.metadata.NoAnyRuleExistsException;
+import 
org.apache.shardingsphere.data.pipeline.core.exception.param.PipelineInvalidParameterException;
 import 
org.apache.shardingsphere.data.pipeline.core.metadata.loader.PipelineSchemaUtil;
 import 
org.apache.shardingsphere.data.pipeline.core.sharding.ShardingColumnsExtractor;
+import 
org.apache.shardingsphere.data.pipeline.core.util.JobDataNodeLineConvertUtil;
 import org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJob;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJobId;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJobType;
@@ -70,6 +69,7 @@ import 
org.apache.shardingsphere.data.pipeline.scenario.migration.config.Migrati
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.context.MigrationProcessContext;
 import org.apache.shardingsphere.data.pipeline.spi.job.JobType;
 import org.apache.shardingsphere.data.pipeline.spi.job.JobTypeFactory;
+import 
org.apache.shardingsphere.data.pipeline.spi.ratelimit.JobRateLimitAlgorithm;
 import 
org.apache.shardingsphere.data.pipeline.spi.sqlbuilder.PipelineSQLBuilder;
 import org.apache.shardingsphere.data.pipeline.util.spi.PipelineTypedSPILoader;
 import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfiguration;
@@ -89,15 +89,17 @@ import 
org.apache.shardingsphere.infra.yaml.config.pojo.YamlRootConfiguration;
 import 
org.apache.shardingsphere.infra.yaml.config.pojo.rule.YamlRuleConfiguration;
 import 
org.apache.shardingsphere.infra.yaml.config.swapper.resource.YamlDataSourceConfigurationSwapper;
 import 
org.apache.shardingsphere.infra.yaml.config.swapper.rule.YamlRuleConfigurationSwapperEngine;
+import 
org.apache.shardingsphere.migration.distsql.statement.MigrateTableStatement;
+import 
org.apache.shardingsphere.migration.distsql.statement.pojo.SourceTargetEntry;
 
 import javax.sql.DataSource;
 import java.nio.charset.StandardCharsets;
 import java.sql.Connection;
-import java.sql.PreparedStatement;
 import java.sql.SQLException;
+import java.sql.Statement;
 import java.util.ArrayList;
 import java.util.Collection;
-import java.util.Collections;
+import java.util.Comparator;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.LinkedHashMap;
@@ -114,17 +116,108 @@ import java.util.stream.Collectors;
 @Slf4j
 public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImpl {
     
+    private static final Gson GSON = new Gson();
+    
     private final YamlRuleConfigurationSwapperEngine swapperEngine = new 
YamlRuleConfigurationSwapperEngine();
     
     private final YamlDataSourceConfigurationSwapper swapper = new 
YamlDataSourceConfigurationSwapper();
     
     private final PipelineDataSourcePersistService dataSourcePersistService = 
new PipelineDataSourcePersistService();
     
+    /**
+     * Create job migration config and start.
+     *
+     * @param param create migration job parameter
+     * @return job id
+     */
+    public String createJobAndStart(final MigrateTableStatement param) {
+        MigrationJobConfiguration jobConfig = new 
YamlMigrationJobConfigurationSwapper().swapToObject(buildYamlJobConfiguration(param));
+        start(jobConfig);
+        return jobConfig.getJobId();
+    }
+    
+    private YamlMigrationJobConfiguration buildYamlJobConfiguration(final 
MigrateTableStatement param) {
+        YamlMigrationJobConfiguration result = new 
YamlMigrationJobConfiguration();
+        result.setTargetDatabaseName(param.getTargetDatabaseName());
+        Map<String, DataSourceProperties> metaDataDataSource = 
dataSourcePersistService.load(new MigrationJobType());
+        Map<String, List<DataNode>> sourceDataNodes = new LinkedHashMap<>();
+        Map<String, YamlPipelineDataSourceConfiguration> configSources = new 
LinkedHashMap<>();
+        List<SourceTargetEntry> sourceTargetEntries = new ArrayList<>(new 
HashSet<>(param.getSourceTargetEntries())).stream().sorted(Comparator.comparing(SourceTargetEntry::getTargetTableName)
+                .thenComparing(each -> 
DataNodeUtil.formatWithSchema(each.getSource()))).collect(Collectors.toList());
+        for (SourceTargetEntry each : sourceTargetEntries) {
+            sourceDataNodes.computeIfAbsent(each.getTargetTableName(), key -> 
new LinkedList<>()).add(each.getSource());
+            ShardingSpherePreconditions.checkState(1 == 
sourceDataNodes.get(each.getTargetTableName()).size(),
+                    () -> new PipelineInvalidParameterException("more than one 
source table for " + each.getTargetTableName()));
+            String dataSourceName = each.getSource().getDataSourceName();
+            if (configSources.containsKey(dataSourceName)) {
+                continue;
+            }
+            Map<String, Object> sourceDataSourceProps = 
swapper.swapToMap(metaDataDataSource.get(dataSourceName));
+            StandardPipelineDataSourceConfiguration sourceDataSourceConfig = 
new StandardPipelineDataSourceConfiguration(sourceDataSourceProps);
+            configSources.put(dataSourceName, 
buildYamlPipelineDataSourceConfiguration(sourceDataSourceConfig.getType(), 
sourceDataSourceConfig.getParameter()));
+            if (null == each.getSource().getSchemaName() && 
sourceDataSourceConfig.getDatabaseType().isSchemaAvailable()) {
+                
each.getSource().setSchemaName(PipelineSchemaUtil.getDefaultSchema(sourceDataSourceConfig));
+            }
+            DatabaseType sourceDatabaseType = 
sourceDataSourceConfig.getDatabaseType();
+            if (null == result.getSourceDatabaseType()) {
+                result.setSourceDatabaseType(sourceDatabaseType.getType());
+            } else if 
(!result.getSourceDatabaseType().equals(sourceDatabaseType.getType())) {
+                throw new PipelineInvalidParameterException("Source storage 
units have different database types");
+            }
+        }
+        result.setSources(configSources);
+        PipelineDataSourceConfiguration targetPipelineDataSourceConfig = 
buildTargetPipelineDataSourceConfiguration(param.getTargetDatabaseName());
+        
result.setTarget(buildYamlPipelineDataSourceConfiguration(targetPipelineDataSourceConfig.getType(),
 targetPipelineDataSourceConfig.getParameter()));
+        
result.setTargetDatabaseType(targetPipelineDataSourceConfig.getDatabaseType().getType());
+        List<JobDataNodeEntry> tablesFirstDataNodes = 
sourceDataNodes.entrySet().stream()
+                .map(entry -> new JobDataNodeEntry(entry.getKey(), 
entry.getValue().subList(0, 1))).collect(Collectors.toList());
+        result.setTargetTableNames(new 
ArrayList<>(sourceDataNodes.keySet()).stream().sorted().collect(Collectors.toList()));
+        
result.setTargetTableSchemaMap(buildTargetTableSchemaMap(sourceDataNodes));
+        result.setTablesFirstDataNodes(new 
JobDataNodeLine(tablesFirstDataNodes).marshal());
+        
result.setJobShardingDataNodes(JobDataNodeLineConvertUtil.convertDataNodesToLines(sourceDataNodes).stream().map(JobDataNodeLine::marshal).collect(Collectors.toList()));
+        extendYamlJobConfiguration(result);
+        return result;
+    }
+    
+    private YamlPipelineDataSourceConfiguration 
buildYamlPipelineDataSourceConfiguration(final String type, final String param) 
{
+        YamlPipelineDataSourceConfiguration result = new 
YamlPipelineDataSourceConfiguration();
+        result.setType(type);
+        result.setParameter(param);
+        return result;
+    }
+    
+    private PipelineDataSourceConfiguration 
buildTargetPipelineDataSourceConfiguration(final String targetDatabaseName) {
+        Map<String, Map<String, Object>> targetDataSourceProps = new 
HashMap<>();
+        ShardingSphereDatabase targetDatabase = 
PipelineContext.getContextManager().getMetaDataContexts().getMetaData().getDatabase(targetDatabaseName);
+        for (Entry<String, DataSource> entry : 
targetDatabase.getResourceMetaData().getDataSources().entrySet()) {
+            Map<String, Object> dataSourceProps = 
swapper.swapToMap(DataSourcePropertiesCreator.create(entry.getValue()));
+            targetDataSourceProps.put(entry.getKey(), dataSourceProps);
+        }
+        YamlRootConfiguration targetRootConfig = 
buildYamlRootConfiguration(targetDatabaseName, targetDataSourceProps, 
targetDatabase.getRuleMetaData().getConfigurations());
+        return new 
ShardingSpherePipelineDataSourceConfiguration(targetRootConfig);
+    }
+    
+    private YamlRootConfiguration buildYamlRootConfiguration(final String 
databaseName, final Map<String, Map<String, Object>> yamlDataSources, final 
Collection<RuleConfiguration> rules) {
+        if (rules.isEmpty()) {
+            throw new NoAnyRuleExistsException(databaseName);
+        }
+        YamlRootConfiguration result = new YamlRootConfiguration();
+        result.setDatabaseName(databaseName);
+        result.setDataSources(yamlDataSources);
+        Collection<YamlRuleConfiguration> yamlRuleConfigurations = 
swapperEngine.swapToYamlRuleConfigurations(rules);
+        result.setRules(yamlRuleConfigurations);
+        return result;
+    }
+    
+    private Map<String, String> buildTargetTableSchemaMap(final Map<String, 
List<DataNode>> sourceDataNodes) {
+        Map<String, String> result = new LinkedHashMap<>();
+        sourceDataNodes.forEach((tableName, dataNodes) -> 
result.put(tableName, dataNodes.get(0).getSchemaName()));
+        return result;
+    }
+    
     @Override
     protected String marshalJobIdLeftPart(final PipelineJobId pipelineJobId) {
-        MigrationJobId jobId = (MigrationJobId) pipelineJobId;
-        String sourceSchemaName = null != jobId.getSourceSchemaName() ? 
jobId.getSourceSchemaName() : "";
-        String text = Joiner.on('|').join(jobId.getSourceResourceName(), 
sourceSchemaName, jobId.getSourceTableName(), jobId.getTargetDatabaseName(), 
jobId.getTargetTableName());
+        String text = GSON.toJson(pipelineJobId);
         return DigestUtils.md5Hex(text.getBytes(StandardCharsets.UTF_8));
     }
     
@@ -133,7 +226,10 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
         JobConfigurationPOJO jobConfigPOJO = getElasticJobConfigPOJO(jobId);
         PipelineJobMetaData jobMetaData = new PipelineJobMetaData(jobId, 
!jobConfigPOJO.isDisabled(),
                 jobConfigPOJO.getShardingTotalCount(), 
jobConfigPOJO.getProps().getProperty("create_time"), 
jobConfigPOJO.getProps().getProperty("stop_time"), 
jobConfigPOJO.getJobParameter());
-        return new TableBasedPipelineJobInfo(jobMetaData, 
getJobConfiguration(jobConfigPOJO).getSourceTableName());
+        List<String> sourceTables = new LinkedList<>();
+        
getJobConfiguration(jobConfigPOJO).getJobShardingDataNodes().forEach(each -> 
each.getEntries().forEach(entry -> entry.getDataNodes()
+                .forEach(dataNode -> 
sourceTables.add(DataNodeUtil.formatWithSchema(dataNode)))));
+        return new TableBasedPipelineJobInfo(jobMetaData, String.join(",", 
sourceTables));
     }
     
     @Override
@@ -142,24 +238,10 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
         if (null == yamlJobConfig.getJobId()) {
             config.setJobId(generateJobId(config));
         }
-        if (Strings.isNullOrEmpty(config.getSourceDatabaseType())) {
-            PipelineDataSourceConfiguration sourceDataSourceConfig = 
PipelineDataSourceConfigurationFactory.newInstance(config.getSource().getType(),
 config.getSource().getParameter());
-            
config.setSourceDatabaseType(sourceDataSourceConfig.getDatabaseType().getType());
-        }
-        if (Strings.isNullOrEmpty(config.getTargetDatabaseType())) {
-            PipelineDataSourceConfiguration targetDataSourceConfig = 
PipelineDataSourceConfigurationFactory.newInstance(config.getTarget().getType(),
 config.getTarget().getParameter());
-            
config.setTargetDatabaseType(targetDataSourceConfig.getDatabaseType().getType());
-        }
-        // target table name is logic table name, source table name is actual 
table name.
-        JobDataNodeEntry nodeEntry = new 
JobDataNodeEntry(config.getTargetTableName(), Collections.singleton(new 
DataNode(config.getSourceResourceName(), config.getSourceTableName())));
-        String dataNodeLine = new 
JobDataNodeLine(Collections.singleton(nodeEntry)).marshal();
-        config.setTablesFirstDataNodes(dataNodeLine);
-        
config.setJobShardingDataNodes(Collections.singletonList(dataNodeLine));
     }
     
     private String generateJobId(final YamlMigrationJobConfiguration config) {
-        MigrationJobId jobId = new 
MigrationJobId(config.getSourceResourceName(), config.getSourceSchemaName(), 
config.getSourceTableName(),
-                config.getTargetDatabaseName(), config.getTargetTableName());
+        MigrationJobId jobId = new 
MigrationJobId(config.getJobShardingDataNodes(), 
config.getTargetDatabaseName());
         return marshalJobId(jobId);
     }
     
@@ -186,25 +268,46 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
     @Override
     public MigrationTaskConfiguration buildTaskConfiguration(final 
PipelineJobConfiguration pipelineJobConfig, final int jobShardingItem, final 
PipelineProcessConfiguration pipelineProcessConfig) {
         MigrationJobConfiguration jobConfig = (MigrationJobConfiguration) 
pipelineJobConfig;
-        Map<ActualTableName, LogicTableName> tableNameMap = new 
LinkedHashMap<>();
-        tableNameMap.put(new ActualTableName(jobConfig.getSourceTableName()), 
new LogicTableName(jobConfig.getTargetTableName()));
-        Map<LogicTableName, String> tableNameSchemaMap = 
TableNameSchemaNameMapping.convert(jobConfig.getSourceSchemaName(), 
Collections.singletonList(jobConfig.getTargetTableName()));
-        TableNameSchemaNameMapping tableNameSchemaNameMapping = new 
TableNameSchemaNameMapping(tableNameSchemaMap);
-        CreateTableConfiguration createTableConfig = 
buildCreateTableConfiguration(jobConfig);
-        DumperConfiguration dumperConfig = 
buildDumperConfiguration(jobConfig.getJobId(), 
jobConfig.getSourceResourceName(), jobConfig.getSource(), tableNameMap, 
tableNameSchemaNameMapping);
+        JobDataNodeLine dataNodeLine = 
jobConfig.getJobShardingDataNodes().get(jobShardingItem);
+        Map<ActualTableName, LogicTableName> tableNameMap = 
buildTableNameMap(dataNodeLine);
+        TableNameSchemaNameMapping tableNameSchemaNameMapping = new 
TableNameSchemaNameMapping(jobConfig.getTargetTableSchemaMap());
+        CreateTableConfiguration createTableConfig = 
buildCreateTableConfiguration(jobConfig, tableNameSchemaNameMapping);
+        String dataSourceName = 
dataNodeLine.getEntries().get(0).getDataNodes().get(0).getDataSourceName();
+        DumperConfiguration dumperConfig = 
buildDumperConfiguration(jobConfig.getJobId(), dataSourceName, 
jobConfig.getSources().get(dataSourceName), tableNameMap, 
tableNameSchemaNameMapping);
+        Set<LogicTableName> targetTableNames = 
jobConfig.getTargetTableNames().stream().map(LogicTableName::new).collect(Collectors.toSet());
         Map<LogicTableName, Set<String>> shardingColumnsMap = new 
ShardingColumnsExtractor().getShardingColumnsMap(
-                ((ShardingSpherePipelineDataSourceConfiguration) 
jobConfig.getTarget()).getRootConfig().getRules(), Collections.singleton(new 
LogicTableName(jobConfig.getTargetTableName())));
+                ((ShardingSpherePipelineDataSourceConfiguration) 
jobConfig.getTarget()).getRootConfig().getRules(), targetTableNames);
         ImporterConfiguration importerConfig = 
buildImporterConfiguration(jobConfig, pipelineProcessConfig, 
shardingColumnsMap, tableNameSchemaNameMapping);
-        return new 
MigrationTaskConfiguration(jobConfig.getSourceResourceName(), 
createTableConfig, dumperConfig, importerConfig);
+        MigrationTaskConfiguration result = new 
MigrationTaskConfiguration(dataSourceName, createTableConfig, dumperConfig, 
importerConfig);
+        log.info("buildTaskConfiguration, result={}", result);
+        return result;
     }
     
-    private CreateTableConfiguration buildCreateTableConfiguration(final 
MigrationJobConfiguration jobConfig) {
-        String sourceSchemaName = jobConfig.getSourceSchemaName();
-        String targetSchemaName = 
PipelineTypedSPILoader.getDatabaseTypedService(DatabaseType.class, 
jobConfig.getTargetDatabaseType()).isSchemaAvailable() ? sourceSchemaName : 
null;
-        CreateTableEntry createTableEntry = new CreateTableEntry(
-                jobConfig.getSource(), new SchemaTableName(new 
SchemaName(sourceSchemaName), new TableName(jobConfig.getSourceTableName())),
-                jobConfig.getTarget(), new SchemaTableName(new 
SchemaName(targetSchemaName), new TableName(jobConfig.getTargetTableName())));
-        return new 
CreateTableConfiguration(Collections.singletonList(createTableEntry));
+    private Map<ActualTableName, LogicTableName> buildTableNameMap(final 
JobDataNodeLine dataNodeLine) {
+        Map<ActualTableName, LogicTableName> result = new LinkedHashMap<>();
+        for (JobDataNodeEntry each : dataNodeLine.getEntries()) {
+            for (DataNode dataNode : each.getDataNodes()) {
+                result.put(new ActualTableName(dataNode.getTableName()), new 
LogicTableName(each.getLogicTableName()));
+            }
+        }
+        return result;
+    }
+    
+    private CreateTableConfiguration buildCreateTableConfiguration(final 
MigrationJobConfiguration jobConfig, final TableNameSchemaNameMapping 
tableNameSchemaNameMapping) {
+        Collection<CreateTableEntry> createTableEntries = new LinkedList<>();
+        for (JobDataNodeEntry each : 
jobConfig.getTablesFirstDataNodes().getEntries()) {
+            String sourceSchemaName = 
tableNameSchemaNameMapping.getSchemaName(each.getLogicTableName());
+            String targetSchemaName = 
TypedSPILoader.getService(DatabaseType.class, 
jobConfig.getTargetDatabaseType()).isSchemaAvailable() ? sourceSchemaName : 
null;
+            DataNode dataNode = each.getDataNodes().get(0);
+            PipelineDataSourceConfiguration sourceDataSourceConfig = 
jobConfig.getSources().get(dataNode.getDataSourceName());
+            CreateTableEntry createTableEntry = new CreateTableEntry(
+                    sourceDataSourceConfig, new SchemaTableName(new 
SchemaName(sourceSchemaName), new TableName(dataNode.getTableName())),
+                    jobConfig.getTarget(), new SchemaTableName(new 
SchemaName(targetSchemaName), new TableName(each.getLogicTableName())));
+            createTableEntries.add(createTableEntry);
+        }
+        CreateTableConfiguration result = new 
CreateTableConfiguration(createTableEntries);
+        log.info("getCreateTableConfiguration, result={}", result);
+        return result;
     }
     
     private DumperConfiguration buildDumperConfiguration(final String jobId, 
final String dataSourceName, final PipelineDataSourceConfiguration 
sourceDataSource,
@@ -220,20 +323,12 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
     
     private ImporterConfiguration buildImporterConfiguration(final 
MigrationJobConfiguration jobConfig, final PipelineProcessConfiguration 
pipelineProcessConfig,
                                                              final 
Map<LogicTableName, Set<String>> shardingColumnsMap, final 
TableNameSchemaNameMapping tableNameSchemaNameMapping) {
+        MigrationProcessContext processContext = new 
MigrationProcessContext(jobConfig.getJobId(), pipelineProcessConfig);
+        JobRateLimitAlgorithm writeRateLimitAlgorithm = 
processContext.getWriteRateLimitAlgorithm();
         int batchSize = pipelineProcessConfig.getWrite().getBatchSize();
         int retryTimes = jobConfig.getRetryTimes();
         int concurrency = jobConfig.getConcurrency();
-        MigrationProcessContext migrationProcessContext = new 
MigrationProcessContext(jobConfig.getJobId(), pipelineProcessConfig);
-        return new ImporterConfiguration(jobConfig.getTarget(), 
unmodifiable(shardingColumnsMap), tableNameSchemaNameMapping, batchSize, 
migrationProcessContext.getWriteRateLimitAlgorithm(),
-                retryTimes, concurrency);
-    }
-    
-    private Map<LogicTableName, Set<String>> unmodifiable(final 
Map<LogicTableName, Set<String>> shardingColumnsMap) {
-        Map<LogicTableName, Set<String>> result = new 
HashMap<>(shardingColumnsMap.size());
-        for (Entry<LogicTableName, Set<String>> entry : 
shardingColumnsMap.entrySet()) {
-            result.put(entry.getKey(), 
Collections.unmodifiableSet(entry.getValue()));
-        }
-        return Collections.unmodifiableMap(result);
+        return new ImporterConfiguration(jobConfig.getTarget(), 
shardingColumnsMap, tableNameSchemaNameMapping, batchSize, 
writeRateLimitAlgorithm, retryTimes, concurrency);
     }
     
     @Override
@@ -305,17 +400,18 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
     
     private void cleanTempTableOnRollback(final String jobId) throws 
SQLException {
         MigrationJobConfiguration jobConfig = getJobConfiguration(jobId);
-        String targetTableName = jobConfig.getTargetTableName();
-        // TODO use jobConfig.targetSchemaName
-        String targetSchemaName = jobConfig.getSourceSchemaName();
         PipelineSQLBuilder pipelineSQLBuilder = 
PipelineTypedSPILoader.getDatabaseTypedService(PipelineSQLBuilder.class, 
jobConfig.getTargetDatabaseType());
+        TableNameSchemaNameMapping mapping = new 
TableNameSchemaNameMapping(jobConfig.getTargetTableSchemaMap());
         try (
                 PipelineDataSourceWrapper dataSource = 
PipelineDataSourceFactory.newInstance(jobConfig.getTarget());
                 Connection connection = dataSource.getConnection()) {
-            String sql = pipelineSQLBuilder.buildDropSQL(targetSchemaName, 
targetTableName);
-            log.info("cleanTempTableOnRollback, targetSchemaName={}, 
targetTableName={}, sql={}", targetSchemaName, targetTableName, sql);
-            try (PreparedStatement preparedStatement = 
connection.prepareStatement(sql)) {
-                preparedStatement.execute();
+            for (String each : jobConfig.getTargetTableNames()) {
+                String targetSchemaName = mapping.getSchemaName(each);
+                String sql = pipelineSQLBuilder.buildDropSQL(targetSchemaName, 
each);
+                log.info("cleanTempTableOnRollback, targetSchemaName={}, 
targetTableName={}, sql={}", targetSchemaName, each, sql);
+                try (Statement statement = connection.createStatement()) {
+                    statement.execute(sql);
+                }
             }
         }
     }
@@ -406,70 +502,6 @@ public final class MigrationJobAPI extends 
AbstractInventoryIncrementalJobAPIImp
         return "";
     }
     
-    /**
-     * Create job migration config and start.
-     *
-     * @param param create migration job parameter
-     * @return job id
-     */
-    public String createJobAndStart(final CreateMigrationJobParameter param) {
-        MigrationJobConfiguration jobConfig = new 
YamlMigrationJobConfigurationSwapper().swapToObject(createYamlJobConfiguration(param));
-        start(jobConfig);
-        return jobConfig.getJobId();
-    }
-    
-    private YamlMigrationJobConfiguration createYamlJobConfiguration(final 
CreateMigrationJobParameter param) {
-        YamlMigrationJobConfiguration result = new 
YamlMigrationJobConfiguration();
-        Map<String, DataSourceProperties> metaDataDataSource = 
dataSourcePersistService.load(new MigrationJobType());
-        Map<String, Object> sourceDataSourceProps = 
swapper.swapToMap(metaDataDataSource.get(param.getSourceResourceName()));
-        StandardPipelineDataSourceConfiguration sourceDataSourceConfig = new 
StandardPipelineDataSourceConfiguration(sourceDataSourceProps);
-        YamlPipelineDataSourceConfiguration sourcePipelineDataSourceConfig = 
createYamlPipelineDataSourceConfiguration(
-                sourceDataSourceConfig.getType(), 
sourceDataSourceConfig.getParameter());
-        result.setSource(sourcePipelineDataSourceConfig);
-        result.setSourceResourceName(param.getSourceResourceName());
-        DatabaseType sourceDatabaseType = 
sourceDataSourceConfig.getDatabaseType();
-        result.setSourceDatabaseType(sourceDatabaseType.getType());
-        String sourceSchemaName = null == param.getSourceSchemaName() && 
sourceDatabaseType.isSchemaAvailable()
-                ? PipelineSchemaUtil.getDefaultSchema(sourceDataSourceConfig)
-                : param.getSourceSchemaName();
-        result.setSourceSchemaName(sourceSchemaName);
-        result.setSourceTableName(param.getSourceTableName());
-        Map<String, Map<String, Object>> targetDataSourceProps = new 
HashMap<>();
-        ShardingSphereDatabase targetDatabase = 
PipelineContext.getContextManager().getMetaDataContexts().getMetaData().getDatabase(param.getTargetDatabaseName());
-        for (Entry<String, DataSource> entry : 
targetDatabase.getResourceMetaData().getDataSources().entrySet()) {
-            Map<String, Object> dataSourceProps = 
swapper.swapToMap(DataSourcePropertiesCreator.create(entry.getValue()));
-            targetDataSourceProps.put(entry.getKey(), dataSourceProps);
-        }
-        String targetDatabaseName = param.getTargetDatabaseName();
-        YamlRootConfiguration targetRootConfig = 
getYamlRootConfiguration(targetDatabaseName, targetDataSourceProps, 
targetDatabase.getRuleMetaData().getConfigurations());
-        PipelineDataSourceConfiguration targetPipelineDataSource = new 
ShardingSpherePipelineDataSourceConfiguration(targetRootConfig);
-        
result.setTarget(createYamlPipelineDataSourceConfiguration(targetPipelineDataSource.getType(),
 targetPipelineDataSource.getParameter()));
-        
result.setTargetDatabaseType(targetPipelineDataSource.getDatabaseType().getType());
-        result.setTargetDatabaseName(targetDatabaseName);
-        result.setTargetTableName(param.getTargetTableName());
-        extendYamlJobConfiguration(result);
-        return result;
-    }
-    
-    private YamlRootConfiguration getYamlRootConfiguration(final String 
databaseName, final Map<String, Map<String, Object>> yamlDataSources, final 
Collection<RuleConfiguration> rules) {
-        if (rules.isEmpty()) {
-            throw new NoAnyRuleExistsException(databaseName);
-        }
-        YamlRootConfiguration result = new YamlRootConfiguration();
-        result.setDatabaseName(databaseName);
-        result.setDataSources(yamlDataSources);
-        Collection<YamlRuleConfiguration> yamlRuleConfigurations = 
swapperEngine.swapToYamlRuleConfigurations(rules);
-        result.setRules(yamlRuleConfigurations);
-        return result;
-    }
-    
-    private YamlPipelineDataSourceConfiguration 
createYamlPipelineDataSourceConfiguration(final String type, final String 
param) {
-        YamlPipelineDataSourceConfiguration result = new 
YamlPipelineDataSourceConfiguration();
-        result.setType(type);
-        result.setParameter(param);
-        return result;
-    }
-    
     @Override
     public JobType getJobType() {
         return JobTypeFactory.getInstance(MigrationJobType.TYPE_CODE);
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 d6254971c18..99f10a02ca4 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
@@ -21,6 +21,8 @@ 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.job.MigrationJobConfiguration;
+import org.apache.shardingsphere.data.pipeline.api.datanode.DataNodeUtil;
+import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceManager;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.job.progress.InventoryIncrementalJobItemProgress;
@@ -32,21 +34,22 @@ import 
org.apache.shardingsphere.data.pipeline.api.metadata.model.PipelineColumn
 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.datasource.DefaultPipelineDataSourceManager;
 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;
 import 
org.apache.shardingsphere.data.pipeline.spi.ratelimit.JobRateLimitAlgorithm;
+import org.apache.shardingsphere.infra.datanode.DataNode;
 import 
org.apache.shardingsphere.infra.util.exception.ShardingSpherePreconditions;
-import 
org.apache.shardingsphere.infra.util.exception.external.sql.type.wrapper.SQLWrapperException;
 
-import java.sql.SQLException;
 import java.util.LinkedHashMap;
+import java.util.LinkedList;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 /**
  * Data consistency checker for migration job.
@@ -69,29 +72,49 @@ public final class MigrationDataConsistencyChecker 
implements PipelineDataConsis
     
     @Override
     public Map<String, DataConsistencyCheckResult> check(final 
DataConsistencyCalculateAlgorithm calculateAlgorithm) {
-        verifyPipelineDatabaseType(calculateAlgorithm, jobConfig.getSource());
+        verifyPipelineDatabaseType(calculateAlgorithm, 
jobConfig.getSources().values().iterator().next());
         verifyPipelineDatabaseType(calculateAlgorithm, jobConfig.getTarget());
-        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());
+        List<String> sourceTableNames = new LinkedList<>();
+        jobConfig.getJobShardingDataNodes().forEach(each -> 
each.getEntries().forEach(entry -> entry.getDataNodes()
+                .forEach(dataNode -> 
sourceTableNames.add(DataNodeUtil.formatWithSchema(dataNode)))));
+        progressContext.setRecordsCount(getRecordsCount());
+        progressContext.getTableNames().addAll(sourceTableNames);
         Map<String, DataConsistencyCheckResult> result = new LinkedHashMap<>();
-        try (
-                PipelineDataSourceWrapper sourceDataSource = 
PipelineDataSourceFactory.newInstance(jobConfig.getSource());
-                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, uniqueKey, metaDataLoader, readRateLimitAlgorithm, 
progressContext);
-            result.put(sourceTable.getTableName().getOriginal(), 
singleTableInventoryChecker.check(calculateAlgorithm));
-        } catch (final SQLException ex) {
-            throw new SQLWrapperException(ex);
+        PipelineDataSourceManager dataSourceManager = new 
DefaultPipelineDataSourceManager();
+        try {
+            AtomicBoolean checkFailed = new AtomicBoolean(false);
+            jobConfig.getJobShardingDataNodes().forEach(each -> 
each.getEntries().forEach(entry -> entry.getDataNodes().forEach(dataNode -> {
+                if (checkFailed.get()) {
+                    return;
+                }
+                DataConsistencyCheckResult checkResult = 
checkSingleTable(entry.getLogicTableName(), dataNode, calculateAlgorithm, 
dataSourceManager);
+                result.put(DataNodeUtil.formatWithSchema(dataNode), 
checkResult);
+                if (!checkResult.isMatched()) {
+                    log.info("unmatched on table '{}', ignore left tables", 
each);
+                    checkFailed.set(true);
+                }
+            })));
+        } finally {
+            dataSourceManager.close();
         }
         return result;
     }
     
+    private DataConsistencyCheckResult checkSingleTable(final String 
targetTableName, final DataNode dataNode,
+                                                        final 
DataConsistencyCalculateAlgorithm calculateAlgorithm, final 
PipelineDataSourceManager dataSourceManager) {
+        SchemaTableName sourceTable = new SchemaTableName(new 
SchemaName(dataNode.getSchemaName()), new TableName(dataNode.getTableName()));
+        SchemaTableName targetTable = new SchemaTableName(new 
SchemaName(dataNode.getSchemaName()), new TableName(targetTableName));
+        PipelineDataSourceWrapper sourceDataSource = 
dataSourceManager.getDataSource(jobConfig.getSources().get(dataNode.getDataSourceName()));
+        PipelineDataSourceWrapper targetDataSource = 
dataSourceManager.getDataSource(jobConfig.getTarget());
+        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, uniqueKey, metaDataLoader, readRateLimitAlgorithm, 
progressContext);
+        return singleTableInventoryChecker.check(calculateAlgorithm);
+    }
+    
     private void verifyPipelineDatabaseType(final 
DataConsistencyCalculateAlgorithm calculateAlgorithm, final 
PipelineDataSourceConfiguration dataSourceConfig) {
         
ShardingSpherePreconditions.checkState(calculateAlgorithm.getSupportedDatabaseTypes().contains(dataSourceConfig.getDatabaseType().getType()),
                 () -> new 
UnsupportedPipelineDatabaseTypeException(dataSourceConfig.getDatabaseType()));
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 0460ac79997..bfde585bee7 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
@@ -23,6 +23,7 @@ import 
org.apache.shardingsphere.data.pipeline.api.config.ingest.InventoryDumper
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceManager;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
+import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
 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.job.progress.JobItemIncrementalTasksProgress;
@@ -49,6 +50,7 @@ import 
org.apache.shardingsphere.mode.lock.GlobalLockDefinition;
 
 import java.sql.SQLException;
 import java.util.Collections;
+import java.util.Map.Entry;
 import java.util.Optional;
 
 /**
@@ -175,10 +177,12 @@ public final class MigrationJobPreparer {
      * @param jobConfig job configuration
      */
     public void cleanup(final MigrationJobConfiguration jobConfig) {
-        try {
-            PipelineJobPreparerUtils.destroyPosition(jobConfig.getJobId(), 
jobConfig.getSource());
-        } catch (final SQLException ex) {
-            log.warn("job destroying failed", ex);
+        for (Entry<String, PipelineDataSourceConfiguration> entry : 
jobConfig.getSources().entrySet()) {
+            try {
+                PipelineJobPreparerUtils.destroyPosition(jobConfig.getJobId(), 
entry.getValue());
+            } catch (final SQLException ex) {
+                log.warn("job destroying failed, jobId={}, dataSourceName={}", 
jobConfig.getJobId(), entry.getKey(), ex);
+            }
         }
     }
 }
diff --git 
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
similarity index 75%
rename from 
kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
rename to 
kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
index eca1818a946..9b478aa8de0 100644
--- 
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
+++ 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfiguration.java
@@ -25,42 +25,35 @@ import 
org.apache.shardingsphere.data.pipeline.api.config.job.yaml.YamlPipelineJ
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.yaml.YamlPipelineDataSourceConfiguration;
 
 import java.util.List;
+import java.util.Map;
 
 /**
  * Migration job configuration for YAML.
  */
 @Getter
 @Setter
-@ToString(exclude = {"source", "target"})
+@ToString(exclude = {"sources", "target"})
 public final class YamlMigrationJobConfiguration implements 
YamlPipelineJobConfiguration {
     
     private String jobId;
     
-    private String sourceResourceName;
-    
     private String targetDatabaseName;
     
-    private String sourceSchemaName;
-    
     private String sourceDatabaseType;
     
     private String targetDatabaseType;
     
-    private String sourceTableName;
-    
-    private String targetTableName;
-    
-    private YamlPipelineDataSourceConfiguration source;
+    private Map<String, YamlPipelineDataSourceConfiguration> sources;
     
     private YamlPipelineDataSourceConfiguration target;
     
+    private List<String> targetTableNames;
+    
     /**
-     * Collection of each logic table's first data node.
-     * <p>
-     * If <pre>actualDataNodes: ds_${0..1}.t_order_${0..1}</pre> and 
<pre>actualDataNodes: ds_${0..1}.t_order_item_${0..1}</pre>,
-     * then value may be: {@code 
t_order:ds_0.t_order_0|t_order_item:ds_0.t_order_item_0}.
-     * </p>
+     * Map{logic table names, schema name}.
      */
+    private Map<String, String> targetTableSchemaMap;
+    
     private String tablesFirstDataNodes;
     
     private List<String> jobShardingDataNodes;
@@ -70,13 +63,13 @@ public final class YamlMigrationJobConfiguration implements 
YamlPipelineJobConfi
     private int retryTimes = 3;
     
     /**
-     * Set source.
+     * Set sources.
      *
-     * @param source source configuration
+     * @param sources source configurations
      */
-    public void setSource(final YamlPipelineDataSourceConfiguration source) {
-        checkParameters(source);
-        this.source = source;
+    public void setSources(final Map<String, 
YamlPipelineDataSourceConfiguration> sources) {
+        sources.values().forEach(this::checkParameters);
+        this.sources = sources;
     }
     
     /**
diff --git 
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
similarity index 72%
rename from 
kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
rename to 
kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
index 601c300abeb..bdd3cb4d788 100644
--- 
a/kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
+++ 
b/kernel/data-pipeline/scenario/migration/src/main/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapper.java
@@ -18,10 +18,15 @@
 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.datanode.JobDataNodeLine;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.yaml.YamlPipelineDataSourceConfigurationSwapper;
 import org.apache.shardingsphere.infra.util.yaml.YamlEngine;
 import 
org.apache.shardingsphere.infra.util.yaml.swapper.YamlConfigurationSwapper;
 
+import java.util.LinkedHashMap;
+import java.util.Map.Entry;
+import java.util.stream.Collectors;
+
 /**
  * YAML migration job configuration swapper.
  */
@@ -33,17 +38,16 @@ public final class YamlMigrationJobConfigurationSwapper 
implements YamlConfigura
     public YamlMigrationJobConfiguration swapToYamlConfiguration(final 
MigrationJobConfiguration data) {
         YamlMigrationJobConfiguration result = new 
YamlMigrationJobConfiguration();
         result.setJobId(data.getJobId());
-        result.setSourceResourceName(data.getSourceResourceName());
         result.setTargetDatabaseName(data.getTargetDatabaseName());
         result.setSourceDatabaseType(data.getSourceDatabaseType());
         result.setTargetDatabaseType(data.getTargetDatabaseType());
-        result.setSourceSchemaName(data.getSourceSchemaName());
-        result.setSourceTableName(data.getSourceTableName());
-        result.setTargetTableName(data.getTargetTableName());
-        
result.setSource(dataSourceConfigSwapper.swapToYamlConfiguration(data.getSource()));
+        
result.setSources(data.getSources().entrySet().stream().collect(Collectors.toMap(Entry::getKey,
+                entry -> 
dataSourceConfigSwapper.swapToYamlConfiguration(entry.getValue()), (key, value) 
-> value, LinkedHashMap::new)));
         
result.setTarget(dataSourceConfigSwapper.swapToYamlConfiguration(data.getTarget()));
-        result.setTablesFirstDataNodes(data.getTablesFirstDataNodes());
-        result.setJobShardingDataNodes(data.getJobShardingDataNodes());
+        result.setTargetTableNames(data.getTargetTableNames());
+        result.setTargetTableSchemaMap(data.getTargetTableSchemaMap());
+        
result.setTablesFirstDataNodes(data.getTablesFirstDataNodes().marshal());
+        
result.setJobShardingDataNodes(data.getJobShardingDataNodes().stream().map(JobDataNodeLine::marshal).collect(Collectors.toList()));
         result.setConcurrency(data.getConcurrency());
         result.setRetryTimes(data.getRetryTimes());
         return result;
@@ -51,12 +55,13 @@ public final class YamlMigrationJobConfigurationSwapper 
implements YamlConfigura
     
     @Override
     public MigrationJobConfiguration swapToObject(final 
YamlMigrationJobConfiguration yamlConfig) {
-        return new MigrationJobConfiguration(yamlConfig.getJobId(), 
yamlConfig.getSourceResourceName(), yamlConfig.getTargetDatabaseName(),
-                yamlConfig.getSourceSchemaName(),
+        return new MigrationJobConfiguration(yamlConfig.getJobId(), 
yamlConfig.getTargetDatabaseName(),
                 yamlConfig.getSourceDatabaseType(), 
yamlConfig.getTargetDatabaseType(),
-                yamlConfig.getSourceTableName(), 
yamlConfig.getTargetTableName(),
-                dataSourceConfigSwapper.swapToObject(yamlConfig.getSource()), 
dataSourceConfigSwapper.swapToObject(yamlConfig.getTarget()),
-                yamlConfig.getTablesFirstDataNodes(), 
yamlConfig.getJobShardingDataNodes(),
+                
yamlConfig.getSources().entrySet().stream().collect(Collectors.toMap(Entry::getKey,
+                        entry -> 
dataSourceConfigSwapper.swapToObject(entry.getValue()), (key, value) -> value, 
LinkedHashMap::new)),
+                dataSourceConfigSwapper.swapToObject(yamlConfig.getTarget()),
+                yamlConfig.getTargetTableNames(), 
yamlConfig.getTargetTableSchemaMap(),
+                
JobDataNodeLine.unmarshal(yamlConfig.getTablesFirstDataNodes()), 
yamlConfig.getJobShardingDataNodes().stream().map(JobDataNodeLine::unmarshal).collect(Collectors.toList()),
                 yamlConfig.getConcurrency(), yamlConfig.getRetryTimes());
     }
     
diff --git 
a/kernel/data-pipeline/scenario/migration/src/test/java/org/apache/shardingsphere/data/pipeline/core/job/PipelineJobIdUtilsTest.java
 
b/kernel/data-pipeline/scenario/migration/src/test/java/org/apache/shardingsphere/data/pipeline/core/job/PipelineJobIdUtilsTest.java
index 2a1d6e902a7..8c0b22a88eb 100644
--- 
a/kernel/data-pipeline/scenario/migration/src/test/java/org/apache/shardingsphere/data/pipeline/core/job/PipelineJobIdUtilsTest.java
+++ 
b/kernel/data-pipeline/scenario/migration/src/test/java/org/apache/shardingsphere/data/pipeline/core/job/PipelineJobIdUtilsTest.java
@@ -22,6 +22,8 @@ import 
org.apache.shardingsphere.data.pipeline.scenario.migration.MigrationJobTy
 import org.apache.shardingsphere.data.pipeline.spi.job.JobType;
 import org.junit.Test;
 
+import java.util.Collections;
+
 import static org.hamcrest.CoreMatchers.instanceOf;
 import static org.hamcrest.MatcherAssert.assertThat;
 
@@ -29,7 +31,7 @@ public final class PipelineJobIdUtilsTest {
     
     @Test
     public void assertParseJobType() {
-        MigrationJobId pipelineJobId = new MigrationJobId("ds_0", null, 
"t_order", "sharding_db", "t_order");
+        MigrationJobId pipelineJobId = new 
MigrationJobId(Collections.singletonList("t_order:ds_0.t_order_0,ds_0.t_order_1"),
 "sharding_db");
         String jobId = 
PipelineJobIdUtils.marshalJobIdCommonPrefix(pipelineJobId) + "abcd";
         JobType actualJobType = PipelineJobIdUtils.parseJobType(jobId);
         assertThat(actualJobType, instanceOf(MigrationJobType.class));
diff --git 
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/MySQLMigrationGeneralE2EIT.java
 
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/MySQLMigrationGeneralE2EIT.java
index 53451df1064..ac8a106f501 100644
--- 
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/MySQLMigrationGeneralE2EIT.java
+++ 
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/MySQLMigrationGeneralE2EIT.java
@@ -95,11 +95,11 @@ public final class MySQLMigrationGeneralE2EIT extends 
AbstractMigrationE2EIT {
         log.info("init data end: {}", LocalDateTime.now());
         startMigration(getSourceTableOrderName(), getTargetTableOrderName());
         startMigration("t_order_item", "t_order_item");
-        String orderJobId = getJobIdByTableName(getSourceTableOrderName());
+        String orderJobId = getJobIdByTableName("ds_0." + 
getSourceTableOrderName());
         waitJobPrepareSuccess(String.format("SHOW MIGRATION STATUS '%s'", 
orderJobId));
         startIncrementTask(new MySQLIncrementTask(getSourceDataSource(), 
getSourceTableOrderName(), new SnowflakeKeyGenerateAlgorithm(), 30));
         assertMigrationSuccessById(orderJobId, "DATA_MATCH");
-        String orderItemJobId = getJobIdByTableName("t_order_item");
+        String orderItemJobId = getJobIdByTableName("ds_0.t_order_item");
         assertMigrationSuccessById(orderItemJobId, "DATA_MATCH");
         ThreadUtil.sleep(2, TimeUnit.SECONDS);
         assertMigrationSuccessById(orderItemJobId, "CRC32_MATCH");
diff --git 
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/PostgreSQLMigrationGeneralE2EIT.java
 
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/PostgreSQLMigrationGeneralE2EIT.java
index 0cc156aad9f..0b953aad941 100644
--- 
a/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/PostgreSQLMigrationGeneralE2EIT.java
+++ 
b/test/e2e/pipeline/src/test/java/org/apache/shardingsphere/test/e2e/data/pipeline/cases/migration/general/PostgreSQLMigrationGeneralE2EIT.java
@@ -109,7 +109,7 @@ public final class PostgreSQLMigrationGeneralE2EIT extends 
AbstractMigrationE2EI
     private void checkOrderMigration() throws SQLException, 
InterruptedException {
         startMigrationWithSchema(getSourceTableOrderName(), "t_order");
         startIncrementTask(new PostgreSQLIncrementTask(getSourceDataSource(), 
PipelineBaseE2EIT.SCHEMA_NAME, getSourceTableOrderName(), 20));
-        String jobId = getJobIdByTableName(getSourceTableOrderName());
+        String jobId = getJobIdByTableName("ds_0.test." + 
getSourceTableOrderName());
         waitIncrementTaskFinished(String.format("SHOW MIGRATION STATUS '%s'", 
jobId));
         stopMigrationByJobId(jobId);
         long recordId = new SnowflakeKeyGenerateAlgorithm().generateKey();
@@ -124,7 +124,7 @@ public final class PostgreSQLMigrationGeneralE2EIT extends 
AbstractMigrationE2EI
     
     private void checkOrderItemMigration() throws SQLException, 
InterruptedException {
         startMigrationWithSchema("t_order_item", "t_order_item");
-        String jobId = getJobIdByTableName("t_order_item");
+        String jobId = getJobIdByTableName("ds_0.test.t_order_item");
         waitIncrementTaskFinished(String.format("SHOW MIGRATION STATUS '%s'", 
jobId));
         assertCheckMigrationSuccess(jobId, "DATA_MATCH");
     }
diff --git a/test/e2e/pipeline/src/test/resources/env/mysql/01-initdb.sql 
b/test/e2e/pipeline/src/test/resources/env/mysql/01-initdb.sql
index 4cb199d48f9..6e04bb35ec1 100644
--- a/test/e2e/pipeline/src/test/resources/env/mysql/01-initdb.sql
+++ b/test/e2e/pipeline/src/test/resources/env/mysql/01-initdb.sql
@@ -17,18 +17,18 @@
 
 REVOKE ALL PRIVILEGES ON *.* FROM 'test_user'@'%';
 
-DROP DATABASE IF EXISTS  pipeline_it_0;
-DROP DATABASE IF EXISTS  pipeline_it_1;
-DROP DATABASE IF EXISTS  pipeline_it_2;
-DROP DATABASE IF EXISTS  pipeline_it_3;
-DROP DATABASE IF EXISTS  pipeline_it_4;
+DROP DATABASE IF EXISTS pipeline_it_0;
+DROP DATABASE IF EXISTS pipeline_it_1;
+DROP DATABASE IF EXISTS pipeline_it_2;
+DROP DATABASE IF EXISTS pipeline_it_3;
+DROP DATABASE IF EXISTS pipeline_it_4;
 CREATE DATABASE pipeline_it_0;
 CREATE DATABASE pipeline_it_1;
 CREATE DATABASE pipeline_it_2;
 CREATE DATABASE pipeline_it_3;
 CREATE DATABASE pipeline_it_4;
 
-GRANT REPLICATION CLIENT, REPLICATION SLAVE ON *.*  TO `test_user`@`%`;
+GRANT REPLICATION CLIENT, REPLICATION SLAVE ON *.* TO `test_user`@`%`;
 -- TODO remove unnecessary permissions
 GRANT CREATE, DROP, SELECT, INSERT, UPDATE, DELETE, INDEX ON pipeline_it_0.* 
TO `test_user`@`%`;
 GRANT CREATE, DROP, SELECT, INSERT, UPDATE, DELETE, INDEX ON pipeline_it_1.* 
TO `test_user`@`%`;
diff --git a/test/e2e/pipeline/src/test/resources/env/postgresql/01-initdb.sql 
b/test/e2e/pipeline/src/test/resources/env/postgresql/01-initdb.sql
index 75035c5463e..0da88254cde 100644
--- a/test/e2e/pipeline/src/test/resources/env/postgresql/01-initdb.sql
+++ b/test/e2e/pipeline/src/test/resources/env/postgresql/01-initdb.sql
@@ -17,11 +17,11 @@
 ALTER USER test_user NOSUPERUSER;
 ALTER USER test_user REPLICATION;
 
-DROP DATABASE IF EXISTS  pipeline_it_0;
-DROP DATABASE IF EXISTS  pipeline_it_1;
-DROP DATABASE IF EXISTS  pipeline_it_2;
-DROP DATABASE IF EXISTS  pipeline_it_3;
-DROP DATABASE IF EXISTS  pipeline_it_4;
+DROP DATABASE IF EXISTS pipeline_it_0;
+DROP DATABASE IF EXISTS pipeline_it_1;
+DROP DATABASE IF EXISTS pipeline_it_2;
+DROP DATABASE IF EXISTS pipeline_it_3;
+DROP DATABASE IF EXISTS pipeline_it_4;
 CREATE DATABASE pipeline_it_0;
 CREATE DATABASE pipeline_it_1;
 CREATE DATABASE pipeline_it_2;
diff --git 
a/test/it/parser/src/main/java/org/apache/shardingsphere/test/it/sql/parser/internal/asserts/statement/ral/impl/migration/update/MigrateTableStatementAssert.java
 
b/test/it/parser/src/main/java/org/apache/shardingsphere/test/it/sql/parser/internal/asserts/statement/ral/impl/migration/update/MigrateTableStatementAssert.java
index 28dec005c37..2a76c27dca4 100644
--- 
a/test/it/parser/src/main/java/org/apache/shardingsphere/test/it/sql/parser/internal/asserts/statement/ral/impl/migration/update/MigrateTableStatementAssert.java
+++ 
b/test/it/parser/src/main/java/org/apache/shardingsphere/test/it/sql/parser/internal/asserts/statement/ral/impl/migration/update/MigrateTableStatementAssert.java
@@ -17,7 +17,9 @@
 
 package 
org.apache.shardingsphere.test.it.sql.parser.internal.asserts.statement.ral.impl.migration.update;
 
+import org.apache.shardingsphere.infra.datanode.DataNode;
 import 
org.apache.shardingsphere.migration.distsql.statement.MigrateTableStatement;
+import 
org.apache.shardingsphere.migration.distsql.statement.pojo.SourceTargetEntry;
 import 
org.apache.shardingsphere.test.it.sql.parser.internal.asserts.SQLCaseAssertContext;
 import 
org.apache.shardingsphere.test.it.sql.parser.internal.cases.parser.jaxb.statement.ral.migration.MigrateTableStatementTestCase;
 
@@ -37,10 +39,13 @@ public final class MigrateTableStatementAssert {
      * @param expected expected migrate table statement test case
      */
     public static void assertIs(final SQLCaseAssertContext assertContext, 
final MigrateTableStatement actual, final MigrateTableStatementTestCase 
expected) {
-        assertThat(assertContext.getText("source database name does not 
match"), actual.getSourceResourceName(), is(expected.getSourceResourceName()));
-        assertThat(assertContext.getText("source schema name does not match"), 
actual.getSourceSchemaName(), is(expected.getSourceSchemaName()));
-        assertThat(assertContext.getText("source table name does not match"), 
actual.getSourceTableName(), is(expected.getSourceTableName()));
         assertThat(assertContext.getText("target database name does not 
match"), actual.getTargetDatabaseName(), is(expected.getTargetDatabaseName()));
-        assertThat(assertContext.getText("target table name does not match"), 
actual.getTargetTableName(), is(expected.getTargetTableName()));
+        assertThat(actual.getSourceTargetEntries().size(), is(1));
+        SourceTargetEntry entry = actual.getSourceTargetEntries().get(0);
+        DataNode dataNode = entry.getSource();
+        assertThat(assertContext.getText("source database name does not 
match"), dataNode.getDataSourceName(), is(expected.getSourceResourceName()));
+        assertThat(assertContext.getText("source schema name does not match"), 
dataNode.getSchemaName(), is(expected.getSourceSchemaName()));
+        assertThat(assertContext.getText("source table name does not match"), 
dataNode.getTableName(), is(expected.getSourceTableName()));
+        assertThat(assertContext.getText("target table name does not match"), 
entry.getTargetTableName(), is(expected.getTargetTableName()));
     }
 }
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
index 94423f7bfa4..ddffef063cb 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/api/datanode/JobDataNodeEntryTest.java
@@ -25,6 +25,7 @@ import java.util.Arrays;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.Assert.assertNotNull;
 
 public final class JobDataNodeEntryTest {
     
@@ -34,4 +35,20 @@ public final class JobDataNodeEntryTest {
         String expected = "t_order:ds_0.t_order_0,ds_0.t_order_1";
         assertThat(actual, is(expected));
     }
+    
+    @Test
+    public void assertUnmarshalWithSchema() {
+        JobDataNodeEntry actual = 
JobDataNodeEntry.unmarshal("t_order:ds_0.public.t_order_0,ds_1.test.t_order_1");
+        assertThat(actual.getLogicTableName(), is("t_order"));
+        assertNotNull(actual.getDataNodes());
+        assertThat(actual.getDataNodes().size(), is(2));
+        DataNode first = actual.getDataNodes().get(0);
+        assertThat(first.getDataSourceName(), is("ds_0"));
+        assertThat(first.getSchemaName(), is("public"));
+        assertThat(first.getTableName(), is("t_order_0"));
+        DataNode second = actual.getDataNodes().get(1);
+        assertThat(second.getDataSourceName(), is("ds_1"));
+        assertThat(second.getSchemaName(), is("test"));
+        assertThat(second.getTableName(), is("t_order_1"));
+    }
 }
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/datasource/DefaultPipelineDataSourceManagerTest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/datasource/DefaultPipelineDataSourceManagerTest.java
index e078e1ab3db..06c0a9f7d5e 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/datasource/DefaultPipelineDataSourceManagerTest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/datasource/DefaultPipelineDataSourceManagerTest.java
@@ -20,6 +20,7 @@ package 
org.apache.shardingsphere.test.it.data.pipeline.core.datasource;
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceManager;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
+import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfigurationFactory;
 import 
org.apache.shardingsphere.data.pipeline.core.datasource.DefaultPipelineDataSourceManager;
 import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
@@ -54,16 +55,17 @@ public final class DefaultPipelineDataSourceManagerTest {
     @Test
     public void assertGetDataSource() {
         PipelineDataSourceManager dataSourceManager = new 
DefaultPipelineDataSourceManager();
-        DataSource actual = dataSourceManager.getDataSource(
-                
PipelineDataSourceConfigurationFactory.newInstance(jobConfig.getSource().getType(),
 jobConfig.getSource().getParameter()));
+        PipelineDataSourceConfiguration source = 
jobConfig.getSources().values().iterator().next();
+        DataSource actual = 
dataSourceManager.getDataSource(PipelineDataSourceConfigurationFactory.newInstance(source.getType(),
 source.getParameter()));
         assertThat(actual, instanceOf(PipelineDataSourceWrapper.class));
     }
     
     @Test
     public void assertClose() throws ReflectiveOperationException {
+        PipelineDataSourceConfiguration source = 
jobConfig.getSources().values().iterator().next();
         PipelineDataSourceManager dataSourceManager = new 
DefaultPipelineDataSourceManager();
         try {
-            
dataSourceManager.getDataSource(PipelineDataSourceConfigurationFactory.newInstance(jobConfig.getSource().getType(),
 jobConfig.getSource().getParameter()));
+            
dataSourceManager.getDataSource(PipelineDataSourceConfigurationFactory.newInstance(source.getType(),
 source.getParameter()));
             
dataSourceManager.getDataSource(PipelineDataSourceConfigurationFactory.newInstance(jobConfig.getTarget().getType(),
 jobConfig.getTarget().getParameter()));
             Map<?, ?> cachedDataSources = (Map<?, ?>) 
Plugins.getMemberAccessor().get(DefaultPipelineDataSourceManager.class.getDeclaredField("cachedDataSources"),
 dataSourceManager);
             assertThat(cachedDataSources.size(), is(2));
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/IncrementalTaskTest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/IncrementalTaskTest.java
index 8871829a359..d7ed050d74a 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/IncrementalTaskTest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/core/task/IncrementalTaskTest.java
@@ -64,7 +64,7 @@ public final class IncrementalTaskTest {
     @Test
     public void assertStart() throws ExecutionException, InterruptedException, 
TimeoutException {
         CompletableFuture.allOf(incrementalTask.start().toArray(new 
CompletableFuture[0])).get(10, TimeUnit.SECONDS);
-        assertThat(incrementalTask.getTaskId(), is("standard_0"));
+        assertThat(incrementalTask.getTaskId(), is("ds_0"));
         assertThat(incrementalTask.getTaskProgress().getPosition(), 
instanceOf(PlaceholderPosition.class));
     }
     
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 73f2d2eb0d7..62df7a72757 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
@@ -32,6 +32,10 @@ import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigur
 import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigurationSwapper;
 import org.apache.shardingsphere.infra.util.spi.type.typed.TypedSPILoader;
 
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
 /**
  * Job configuration builder.
  */
@@ -39,28 +43,43 @@ import 
org.apache.shardingsphere.infra.util.spi.type.typed.TypedSPILoader;
 public final class JobConfigurationBuilder {
     
     /**
-     * Create job configuration.
+     * Create migration job configuration.
      *
      * @return created job configuration
      */
+    // TODO Rename createJobConfiguration
     public static MigrationJobConfiguration createJobConfiguration() {
+        return new 
YamlMigrationJobConfigurationSwapper().swapToObject(createYamlMigrationJobConfiguration());
+    }
+    
+    /**
+     * Create YAML migration job configuration.
+     *
+     * @return created job configuration
+     */
+    public static YamlMigrationJobConfiguration 
createYamlMigrationJobConfiguration() {
         YamlMigrationJobConfiguration result = new 
YamlMigrationJobConfiguration();
         result.setTargetDatabaseName("logic_db");
-        result.setSourceResourceName("standard_0");
-        result.setSourceTableName("t_order");
-        result.setTargetTableName("t_order");
+        result.setSourceDatabaseType("H2");
+        result.setTargetDatabaseType("H2");
+        result.setTargetTableNames(Collections.singletonList("t_order"));
+        Map<String, String> targetTableSchemaMap = new LinkedHashMap<>();
+        targetTableSchemaMap.put("t_order", "");
+        result.setTargetTableSchemaMap(targetTableSchemaMap);
+        result.setTablesFirstDataNodes("t_order:ds_0.t_order");
+        
result.setJobShardingDataNodes(Collections.singletonList("t_order:ds_0.t_order"));
         result.setJobId(generateJobId(result));
-        result.setSource(createYamlPipelineDataSourceConfiguration(new 
StandardPipelineDataSourceConfiguration(ConfigurationFileUtil.readFile("migration_standard_jdbc_source.yaml"))));
+        Map<String, YamlPipelineDataSourceConfiguration> sources = new 
LinkedHashMap<>();
+        sources.put("ds_0", createYamlPipelineDataSourceConfiguration(new 
StandardPipelineDataSourceConfiguration(ConfigurationFileUtil.readFile("migration_standard_jdbc_source.yaml"))));
+        result.setSources(sources);
         result.setTarget(createYamlPipelineDataSourceConfiguration(new 
ShardingSpherePipelineDataSourceConfiguration(
                 
ConfigurationFileUtil.readFile("migration_sharding_sphere_jdbc_target.yaml"))));
         TypedSPILoader.getService(PipelineJobAPI.class, 
"MIGRATION").extendYamlJobConfiguration(result);
-        return new YamlMigrationJobConfigurationSwapper().swapToObject(result);
+        return result;
     }
     
     private static String generateJobId(final YamlMigrationJobConfiguration 
yamlJobConfig) {
-        String sourceTableName = RandomStringUtils.randomAlphabetic(32);
-        MigrationJobId migrationJobId = new 
MigrationJobId(yamlJobConfig.getSourceResourceName(), 
yamlJobConfig.getSourceSchemaName(), sourceTableName,
-                yamlJobConfig.getTargetDatabaseName(), 
yamlJobConfig.getTargetTableName());
+        MigrationJobId migrationJobId = new 
MigrationJobId(yamlJobConfig.getJobShardingDataNodes(), 
RandomStringUtils.randomAlphabetic(32));
         return new MigrationJobAPI().marshalJobId(migrationJobId);
     }
     
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/api/impl/MigrationJobAPITest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/api/impl/MigrationJobAPITest.java
index fb1994ee3ed..3375cc519ea 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/api/impl/MigrationJobAPITest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/api/impl/MigrationJobAPITest.java
@@ -22,12 +22,13 @@ import 
org.apache.shardingsphere.data.pipeline.api.check.consistency.DataConsist
 import 
org.apache.shardingsphere.data.pipeline.api.check.consistency.DataConsistencyContentCheckResult;
 import 
org.apache.shardingsphere.data.pipeline.api.check.consistency.DataConsistencyCountCheckResult;
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
+import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeEntry;
+import org.apache.shardingsphere.data.pipeline.api.datanode.JobDataNodeLine;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.PipelineDataSourceWrapper;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfiguration;
 import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDataSourceConfigurationFactory;
 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.pojo.CreateMigrationJobParameter;
 import 
org.apache.shardingsphere.data.pipeline.api.pojo.InventoryIncrementalJobItemInfo;
 import org.apache.shardingsphere.data.pipeline.core.api.PipelineAPIFactory;
 import 
org.apache.shardingsphere.data.pipeline.core.api.impl.PipelineDataSourcePersistService;
@@ -42,10 +43,13 @@ import 
org.apache.shardingsphere.data.pipeline.spi.datasource.creator.PipelineDa
 import org.apache.shardingsphere.elasticjob.infra.pojo.JobConfigurationPOJO;
 import org.apache.shardingsphere.infra.database.type.DatabaseType;
 import org.apache.shardingsphere.infra.database.type.DatabaseTypeEngine;
+import org.apache.shardingsphere.infra.datanode.DataNode;
 import 
org.apache.shardingsphere.infra.datasource.pool.creator.DataSourcePoolCreator;
 import org.apache.shardingsphere.infra.datasource.props.DataSourceProperties;
 import org.apache.shardingsphere.infra.util.spi.type.typed.TypedSPILoader;
 import org.apache.shardingsphere.infra.util.yaml.YamlEngine;
+import 
org.apache.shardingsphere.migration.distsql.statement.MigrateTableStatement;
+import 
org.apache.shardingsphere.migration.distsql.statement.pojo.SourceTargetEntry;
 import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
 import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.PipelineContextUtil;
 import org.junit.AfterClass;
@@ -178,9 +182,10 @@ public final class MigrationJobAPITest {
         DataConsistencyCalculateAlgorithm calculateAlgorithm = 
jobAPI.buildDataConsistencyCalculateAlgorithm(jobConfig, "FIXTURE", null);
         Map<String, DataConsistencyCheckResult> checkResultMap = 
jobAPI.dataConsistencyCheck(jobConfig, calculateAlgorithm, new 
ConsistencyCheckJobItemProgressContext(jobId.get(), 0));
         assertThat(checkResultMap.size(), is(1));
-        
assertTrue(checkResultMap.get("t_order").getCountCheckResult().isMatched());
-        
assertThat(checkResultMap.get("t_order").getCountCheckResult().getTargetRecordsCount(),
 is(2L));
-        
assertTrue(checkResultMap.get("t_order").getContentCheckResult().isMatched());
+        String checkKey = "ds_0.t_order";
+        
assertTrue(checkResultMap.get(checkKey).getCountCheckResult().isMatched());
+        
assertThat(checkResultMap.get(checkKey).getCountCheckResult().getTargetRecordsCount(),
 is(2L));
+        
assertTrue(checkResultMap.get(checkKey).getContentCheckResult().isMatched());
     }
     
     @Test
@@ -236,7 +241,8 @@ public final class MigrationJobAPITest {
     
     @SneakyThrows(SQLException.class)
     private void initTableData(final MigrationJobConfiguration jobConfig) {
-        PipelineDataSourceConfiguration sourceDataSourceConfig = 
PipelineDataSourceConfigurationFactory.newInstance(jobConfig.getSource().getType(),
 jobConfig.getSource().getParameter());
+        PipelineDataSourceConfiguration source = 
jobConfig.getSources().values().iterator().next();
+        PipelineDataSourceConfiguration sourceDataSourceConfig = 
PipelineDataSourceConfigurationFactory.newInstance(source.getType(), 
source.getParameter());
         initTableData(TypedSPILoader.getService(
                 PipelineDataSourceCreator.class, 
sourceDataSourceConfig.getType()).createPipelineDataSource(sourceDataSourceConfig.getDataSourceConfiguration()));
         PipelineDataSourceConfiguration targetDataSourceConfig = 
PipelineDataSourceConfigurationFactory.newInstance(jobConfig.getTarget().getType(),
 jobConfig.getTarget().getParameter());
@@ -275,12 +281,19 @@ public final class MigrationJobAPITest {
     @Test
     public void assertCreateJobConfig() throws SQLException {
         initIntPrimaryEnvironment();
-        String jobId = jobAPI.createJobAndStart(new 
CreateMigrationJobParameter("ds_0", null, "t_order", "logic_db", "t_order"));
-        MigrationJobConfiguration jobConfig = 
jobAPI.getJobConfiguration(jobId);
-        assertThat(jobConfig.getSourceResourceName(), is("ds_0"));
-        assertThat(jobConfig.getSourceTableName(), is("t_order"));
-        assertThat(jobConfig.getTargetDatabaseName(), is("logic_db"));
-        assertThat(jobConfig.getTargetTableName(), is("t_order"));
+        SourceTargetEntry sourceTargetEntry = new 
SourceTargetEntry("logic_db", new DataNode("ds_0", "t_order"), "t_order");
+        String jobId = jobAPI.createJobAndStart(new 
MigrateTableStatement(Collections.singletonList(sourceTargetEntry), 
"logic_db"));
+        MigrationJobConfiguration actual = jobAPI.getJobConfiguration(jobId);
+        assertThat(actual.getTargetDatabaseName(), is("logic_db"));
+        List<JobDataNodeLine> dataNodeLines = actual.getJobShardingDataNodes();
+        assertThat(dataNodeLines.size(), is(1));
+        assertThat(dataNodeLines.get(0).getEntries().size(), is(1));
+        JobDataNodeEntry entry = dataNodeLines.get(0).getEntries().get(0);
+        assertThat(entry.getDataNodes().size(), is(1));
+        DataNode dataNode = entry.getDataNodes().get(0);
+        assertThat(dataNode.getDataSourceName(), is("ds_0"));
+        assertThat(dataNode.getTableName(), is("t_order"));
+        assertThat(entry.getLogicTableName(), is("t_order"));
     }
     
     private void initIntPrimaryEnvironment() throws SQLException {
diff --git 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyCheckerTest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyCheckerTest.java
index aa46cc5ec52..8065d0e8c92 100644
--- 
a/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyCheckerTest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/check/consistency/MigrationDataConsistencyCheckerTest.java
@@ -23,15 +23,15 @@ import 
org.apache.shardingsphere.data.pipeline.api.datasource.config.PipelineDat
 import org.apache.shardingsphere.data.pipeline.core.api.PipelineAPIFactory;
 import 
org.apache.shardingsphere.data.pipeline.core.check.consistency.ConsistencyCheckJobItemProgressContext;
 import 
org.apache.shardingsphere.data.pipeline.core.datasource.DefaultPipelineDataSourceManager;
-import 
org.apache.shardingsphere.test.it.data.pipeline.core.fixture.DataConsistencyCalculateAlgorithmFixture;
-import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
-import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.PipelineContextUtil;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.check.consistency.MigrationDataConsistencyChecker;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.context.MigrationJobItemContext;
 import 
org.apache.shardingsphere.data.pipeline.scenario.migration.context.MigrationProcessContext;
 import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigurationSwapper;
 import org.apache.shardingsphere.elasticjob.infra.pojo.JobConfigurationPOJO;
 import org.apache.shardingsphere.infra.util.yaml.YamlEngine;
+import 
org.apache.shardingsphere.test.it.data.pipeline.core.fixture.DataConsistencyCalculateAlgorithmFixture;
+import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
+import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.PipelineContextUtil;
 import org.junit.BeforeClass;
 import org.junit.Test;
 
@@ -60,9 +60,10 @@ public final class MigrationDataConsistencyCheckerTest {
         
PipelineAPIFactory.getGovernanceRepositoryAPI().persistJobItemProgress(jobConfig.getJobId(),
 0, "");
         Map<String, DataConsistencyCheckResult> actual = new 
MigrationDataConsistencyChecker(jobConfig, new 
MigrationProcessContext(jobConfig.getJobId(), null),
                 createConsistencyCheckJobItemProgressContext()).check(new 
DataConsistencyCalculateAlgorithmFixture());
-        assertTrue(actual.get("t_order").getCountCheckResult().isMatched());
-        
assertThat(actual.get("t_order").getCountCheckResult().getSourceRecordsCount(), 
is(actual.get("t_order").getCountCheckResult().getTargetRecordsCount()));
-        assertTrue(actual.get("t_order").getContentCheckResult().isMatched());
+        String checkKey = "ds_0.t_order";
+        assertTrue(actual.get(checkKey).getCountCheckResult().isMatched());
+        
assertThat(actual.get(checkKey).getCountCheckResult().getSourceRecordsCount(), 
is(actual.get(checkKey).getCountCheckResult().getTargetRecordsCount()));
+        assertTrue(actual.get(checkKey).getContentCheckResult().isMatched());
     }
     
     private ConsistencyCheckJobItemProgressContext 
createConsistencyCheckJobItemProgressContext() {
diff --git 
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/yaml/YamlMigrationJobConfigurationSwapperTest.java
similarity index 54%
rename from 
kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
rename to 
test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/yaml/YamlMigrationJobConfigurationSwapperTest.java
index 05e0f6b21e6..59860398666 100644
--- 
a/kernel/data-pipeline/core/src/test/java/org/apache/shardingsphere/data/pipeline/yaml/job/YamlMigrationJobConfigurationSwapperTest.java
+++ 
b/test/it/pipeline/src/test/java/org/apache/shardingsphere/test/it/data/pipeline/scenario/migration/yaml/YamlMigrationJobConfigurationSwapperTest.java
@@ -15,9 +15,13 @@
  * limitations under the License.
  */
 
-package org.apache.shardingsphere.data.pipeline.yaml.job;
+package 
org.apache.shardingsphere.test.it.data.pipeline.scenario.migration.yaml;
 
 import 
org.apache.shardingsphere.data.pipeline.api.config.job.MigrationJobConfiguration;
+import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfiguration;
+import 
org.apache.shardingsphere.data.pipeline.yaml.job.YamlMigrationJobConfigurationSwapper;
+import org.apache.shardingsphere.infra.util.yaml.YamlEngine;
+import 
org.apache.shardingsphere.test.it.data.pipeline.core.util.JobConfigurationBuilder;
 import org.junit.Test;
 
 import static org.hamcrest.CoreMatchers.is;
@@ -26,9 +30,11 @@ import static org.hamcrest.MatcherAssert.assertThat;
 public final class YamlMigrationJobConfigurationSwapperTest {
     
     @Test
-    public void assertSwapToObject() {
-        YamlMigrationJobConfiguration yamlJobConfig = new 
YamlMigrationJobConfiguration();
-        MigrationJobConfiguration actual = new 
YamlMigrationJobConfigurationSwapper().swapToObject(yamlJobConfig);
-        assertThat(actual.getJobShardingCount(), is(1));
+    public void assertMarsharlUnmarshal() {
+        YamlMigrationJobConfiguration yamlJobConfig = 
JobConfigurationBuilder.createYamlMigrationJobConfiguration();
+        YamlMigrationJobConfigurationSwapper swapper = new 
YamlMigrationJobConfigurationSwapper();
+        MigrationJobConfiguration jobConfig = 
swapper.swapToObject(yamlJobConfig);
+        YamlMigrationJobConfiguration actual = 
swapper.swapToYamlConfiguration(jobConfig);
+        assertThat(YamlEngine.marshal(actual), 
is(YamlEngine.marshal(yamlJobConfig)));
     }
 }

Reply via email to