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)));
}
}