Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion dt-common/src/config/checker_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ pub const DEFAULT_CDC_CHECK_LOG_INTERVAL_SECS: u64 = 30;
/// Common checker settings.
///
/// Standalone snapshot/struct/check-log tasks use `[sinker] sink_type=check`; the checker target
/// connection is loaded through the regular MySQL/PostgreSQL/MongoDB sinker configuration.
/// connection is loaded through the regular MySQL/PostgreSQL/MSSQL/MongoDB sinker configuration.
/// `[checker_output]` owns result output settings. CDC inline check is enabled separately through
/// `[checker_cdc] is_enabled=true`.
#[derive(Clone, Debug)]
Expand Down
7 changes: 7 additions & 0 deletions dt-common/src/config/extractor_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,13 @@ pub enum ExtractorConfig {
db_batch_size: usize,
},

MssqlStruct {
url: String,
connection_auth: ConnectionAuthConfig,
dbs: Vec<String>,
db_batch_size: usize,
},

MysqlSnapshot {
url: String,
connection_auth: ConnectionAuthConfig,
Expand Down
6 changes: 6 additions & 0 deletions dt-common/src/config/sinker_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,12 @@ pub enum SinkerConfig {
conflict_policy: ConflictPolicyEnum,
},

MssqlStruct {
url: String,
connection_auth: ConnectionAuthConfig,
conflict_policy: ConflictPolicyEnum,
},

Kafka {
url: String,
batch_size: usize,
Expand Down
150 changes: 146 additions & 4 deletions dt-common/src/config/task_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,7 +296,7 @@ impl TaskConfig {
}
} else {
match (kind, sink_type, target_db_type) {
(TaskKind::Struct, SinkType::Check, DbType::Mysql | DbType::Pg) => {
(TaskKind::Struct, SinkType::Check, DbType::Mysql | DbType::Pg | DbType::Mssql) => {
Some(CheckMode::Standalone)
}
(
Expand Down Expand Up @@ -543,6 +543,25 @@ impl TaskConfig {
partition_cols: loader.get_optional(EXTRACTOR, PARTITION_COLS)?,
}
}
ExtractType::Struct => {
Self::validate_mssql_connection(
EXTRACTOR,
&url,
&connection_auth,
basic.app_name.as_deref(),
max_connections,
)?;
ExtractorConfig::MssqlStruct {
url,
connection_auth,
dbs: Vec::new(),
db_batch_size: loader.get_with_default(
EXTRACTOR,
"db_batch_size",
DEFAULT_DB_BATCH_SIZE,
)?,
}
}
_ => bail! { not_supported_err },
},

Expand Down Expand Up @@ -811,7 +830,7 @@ impl TaskConfig {
},

DbType::Mssql => match sink_type {
SinkType::Write => {
SinkType::Write | SinkType::Check => {
Self::validate_mssql_connection(
SINKER,
&url,
Expand All @@ -826,6 +845,20 @@ impl TaskConfig {
replace: loader.get_with_default(SINKER, REPLACE, true)?,
}
}
SinkType::Struct => {
Self::validate_mssql_connection(
SINKER,
&url,
&connection_auth,
basic.app_name.as_deref(),
max_connections,
)?;
SinkerConfig::MssqlStruct {
url,
connection_auth,
conflict_policy,
}
}
_ => bail! { not_supported_err },
},

Expand Down Expand Up @@ -1475,8 +1508,8 @@ mod tests {
};

use super::{
CheckMode, DbType, ExtractorConfig, ParallelType, RdbParallelType, SinkerConfig,
TaskConfig, TaskKind, TaskType,
CheckMode, ConflictPolicyEnum, DbType, ExtractorConfig, ParallelType, RdbParallelType,
SinkerConfig, TaskConfig, TaskKind, TaskType,
};
use crate::config::parallelizer_config::{
ChunkPartitionerRebalanceCost, ChunkPartitionerRebalanceStrategy,
Expand Down Expand Up @@ -2158,6 +2191,115 @@ parallel_size=2
));
}

#[test]
fn mssql_struct_config_is_loaded() {
let config = load_temp_task_config(
r#"[extractor]
db_type=mssql
extract_type=struct
url=server=tcp:127.0.0.1,1433;database=ape_dts
username=sa
password=Password123!
ssl_mode=disable
db_batch_size=11

[sinker]
db_type=mssql
sink_type=struct
url=server=tcp:127.0.0.1,1434;database=ape_dts
username=sa
password=Password123!
ssl_mode=disable
conflict_policy=ignore
"#,
)
.expect("MSSQL struct config should be accepted");

assert_eq!(
config.task_type(),
Some(TaskType::new(TaskKind::Struct, None))
);
assert!(matches!(
config.extractor,
ExtractorConfig::MssqlStruct {
db_batch_size: 11,
..
}
));
assert!(matches!(
config.sinker,
SinkerConfig::MssqlStruct {
conflict_policy: ConflictPolicyEnum::Ignore,
..
}
));
}

#[test]
fn mssql_struct_check_config_is_loaded() {
let config = load_temp_task_config(
r#"[extractor]
db_type=mssql
extract_type=struct
url=server=tcp:127.0.0.1,1433;database=ape_dts
username=sa
password=Password123!
ssl_mode=disable

[sinker]
db_type=mssql
sink_type=check
url=server=tcp:127.0.0.1,1434;database=ape_dts
username=sa
password=Password123!
ssl_mode=disable
app_name=mssql-struct-check
connection_timeout_secs=9
max_connections=3
"#,
)
.expect("MSSQL standalone struct check config should be accepted");

assert_eq!(
config.task_type(),
Some(TaskType::new(TaskKind::Struct, Some(CheckMode::Standalone)))
);
assert_eq!(
config.checker_target().expect("checker target").db_type,
DbType::Mssql
);
assert!(config.checker.is_some());
assert!(matches!(config.sinker, SinkerConfig::Mssql { .. }));
}

#[test]
fn mssql_snapshot_check_config_remains_unsupported() {
let result = load_temp_task_config(
r#"[extractor]
db_type=mssql
extract_type=snapshot
url=sqlserver://127.0.0.1:1433?database=ape_dts
username=sa
password=Password123!
ssl_mode=disable

[sinker]
db_type=mssql
sink_type=check
url=server=tcp:127.0.0.1,1434;database=ape_dts
username=sa
password=Password123!
ssl_mode=disable

[parallelizer]
parallel_type=snapshot
parallel_size=1
"#,
);

assert!(result.is_err());
}

#[test]
fn mssql_snapshot_config_defaults_connection_timeout_to_15_seconds() {
let config = load_temp_task_config(
Expand Down
72 changes: 60 additions & 12 deletions dt-common/src/meta/ddl_meta/ddl_data.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,29 +27,77 @@ impl DdlData {
}

pub fn get_schema_tb(&self) -> (String, String) {
let (mut schema, tb) = self.statement.get_schema_tb();
if schema.is_empty() {
schema = self.default_schema.clone()
let (db, schema, tb) = self.get_db_schema_tb();
match self.ddl_type {
DdlType::CreateDatabase | DdlType::DropDatabase | DdlType::AlterDatabase => {
(db, String::new())
}
DdlType::CreateSchema | DdlType::DropSchema | DdlType::AlterSchema => {
(schema, String::new())
}
_ => (schema, tb),
}
(schema, tb)
}

pub fn get_db_schema_tb(&self) -> (String, String, String) {
let (schema, tb) = self.get_schema_tb();
(self.default_db.clone(), schema, tb)
}
let (mut db, mut schema, tb) = self.statement.get_db_schema_tb();
match self.ddl_type {
DdlType::CreateDatabase | DdlType::DropDatabase | DdlType::AlterDatabase => {
return (db, String::new(), String::new());
}
DdlType::CreateSchema | DdlType::DropSchema | DdlType::AlterSchema => {
if db.is_empty() {
db = self.default_db.clone();
}
return (db, schema, String::new());
}
_ => {}
}

pub fn get_rename_to_schema_tb(&self) -> (String, String) {
let (mut schema, tb) = self.statement.get_rename_to_schema_tb();
if db.is_empty() {
db = self.default_db.clone();
}
if schema.is_empty() {
schema = self.default_schema.clone()
schema = self.default_schema.clone();
}
(db, schema, tb)
}

pub fn get_rename_to_schema_tb(&self) -> (String, String) {
let (_, schema, tb) = self.get_rename_to_db_schema_tb();
(schema, tb)
}

pub fn get_rename_to_db_schema_tb(&self) -> (String, String, String) {
let (schema, tb) = self.get_rename_to_schema_tb();
(self.default_db.clone(), schema, tb)
let (mut db, mut schema, tb) = self.statement.get_rename_to_db_schema_tb();
if tb.is_empty() {
return (String::new(), String::new(), String::new());
}

let (src_db, src_schema, _) = self.get_db_schema_tb();
if db.is_empty() {
db = src_db;
}
if schema.is_empty() {
schema = src_schema;
}
(db, schema, tb)
}

pub fn route(&mut self, dst_db: String, dst_schema: String, dst_tb: String) {
self.statement
.route_db_schema_tb(dst_db, dst_schema, dst_tb);
}

pub fn route_rename(
&mut self,
dst_schema: String,
dst_tb: String,
dst_new_schema: String,
dst_new_tb: String,
) {
self.statement
.route_rename(dst_schema, dst_tb, dst_new_schema, dst_new_tb);
}

pub fn split_to_multi(self) -> Vec<DdlData> {
Expand Down
Loading
Loading