1use std::collections::HashMap;
16
17use anyhow::{Context, anyhow};
18use chrono::{DateTime, NaiveDateTime};
19use futures::stream::BoxStream;
20use futures::{StreamExt, pin_mut, stream};
21use futures_async_stream::try_stream;
22use itertools::Itertools;
23use mysql_async::prelude::*;
24use mysql_common::params::Params;
25use mysql_common::value::Value;
26use risingwave_common::bail;
27use risingwave_common::catalog::{CDC_OFFSET_COLUMN_NAME, ColumnDesc, ColumnId, Field, Schema};
28use risingwave_common::row::OwnedRow;
29use risingwave_common::types::{DataType, Datum, Decimal, F32, ScalarImpl};
30use risingwave_common::util::iter_util::ZipEqFast;
31use sea_schema::mysql::def::{ColumnDefault, ColumnKey, ColumnType};
32use sea_schema::mysql::discovery::SchemaDiscovery;
33use sea_schema::mysql::query::SchemaQueryBuilder;
34use sea_schema::sea_query::{Alias, IntoIden};
35use serde::{Deserialize, Serialize};
36use sqlx::MySqlPool;
37use sqlx::mysql::MySqlConnectOptions;
38use thiserror_ext::AsReport;
39
40use crate::connector_common::SslMode;
41pub use crate::connector_common::SslMode as MySqlSslMode;
43use crate::error::{ConnectorError, ConnectorResult};
44use crate::source::CdcTableSnapshotSplit;
45use crate::source::cdc::external::{
46 CdcOffset, CdcOffsetParseFunc, CdcTableSnapshotSplitOption, DebeziumOffset,
47 ExternalTableConfig, ExternalTableReader, SchemaTableName, mysql_row_to_owned_row,
48};
49
50pub fn build_mysql_connection_pool(
67 host: &str,
68 port: u16,
69 username: &str,
70 password: &str,
71 database: &str,
72 ssl_mode: SslMode,
73) -> mysql_async::Pool {
74 let mut opts_builder = mysql_async::OptsBuilder::default()
75 .user(Some(username))
76 .pass(Some(password))
77 .ip_or_hostname(host)
78 .tcp_port(port)
79 .db_name(Some(database));
80
81 opts_builder = match ssl_mode {
82 SslMode::Disabled | SslMode::Preferred => opts_builder.ssl_opts(None),
83 SslMode::Required | SslMode::VerifyCa | SslMode::VerifyFull => {
85 let ssl_without_verify = mysql_async::SslOpts::default()
86 .with_danger_accept_invalid_certs(true)
87 .with_danger_skip_domain_validation(true);
88 opts_builder.ssl_opts(Some(ssl_without_verify))
89 }
90 };
91
92 mysql_async::Pool::new(opts_builder)
93}
94
95#[derive(Debug, Clone, Default, PartialEq, PartialOrd, Serialize, Deserialize)]
96pub struct MySqlOffset {
97 pub filename: String,
98 pub position: u64,
99}
100
101impl MySqlOffset {
102 pub fn new(filename: String, position: u64) -> Self {
103 Self { filename, position }
104 }
105}
106
107impl MySqlOffset {
108 pub fn parse_debezium_offset(offset: &str) -> ConnectorResult<Self> {
109 let dbz_offset: DebeziumOffset = serde_json::from_str(offset)
110 .with_context(|| format!("invalid upstream offset: {}", offset))?;
111
112 Ok(Self {
113 filename: dbz_offset
114 .source_offset
115 .file
116 .context("binlog file not found in offset")?,
117 position: dbz_offset
118 .source_offset
119 .pos
120 .context("binlog position not found in offset")?,
121 })
122 }
123}
124
125pub struct MySqlExternalTable {
126 column_descs: Vec<ColumnDesc>,
127 pk_names: Vec<String>,
128}
129
130impl MySqlExternalTable {
131 pub async fn connect(config: ExternalTableConfig) -> ConnectorResult<Self> {
132 tracing::debug!("connect to mysql");
133 let options = MySqlConnectOptions::new()
134 .username(&config.username)
135 .password(&config.password)
136 .host(&config.host)
137 .port(config.port.parse::<u16>().unwrap())
138 .database(&config.database)
139 .ssl_mode(match config.ssl_mode {
140 SslMode::Disabled => sqlx::mysql::MySqlSslMode::Disabled,
141 SslMode::Preferred => sqlx::mysql::MySqlSslMode::Preferred,
142 SslMode::Required => sqlx::mysql::MySqlSslMode::Required,
143 _ => {
144 return Err(anyhow!("unsupported SSL mode").into());
145 }
146 });
147
148 let connection = MySqlPool::connect_with(options).await?;
149 let mut schema_discovery = SchemaDiscovery::new(connection, config.database.as_str());
150
151 let system_info = schema_discovery.discover_system().await?;
153 schema_discovery.query = SchemaQueryBuilder::new(system_info.clone());
154 let schema = Alias::new(config.database.as_str()).into_iden();
155 let table = Alias::new(config.table.as_str()).into_iden();
156 let columns = schema_discovery
157 .discover_columns(schema, table, &system_info)
158 .await?;
159 let mut column_descs = vec![];
160 let mut pk_names = vec![];
161 for col in columns {
162 let data_type = mysql_type_to_rw_type(&col.col_type)?;
163 let col_name = col.name.to_lowercase();
165 let column_desc = if let Some(default) = col.default {
166 let snapshot_value = derive_default_value(default.clone(), &data_type)
167 .unwrap_or_else(|e| {
168 tracing::warn!(
169 column = col_name,
170 ?default,
171 %data_type,
172 error = %e.as_report(),
173 "failed to derive column default value, fallback to `NULL`",
174 );
175 None
176 });
177
178 ColumnDesc::named_with_default_value(
179 col_name.clone(),
180 ColumnId::placeholder(),
181 data_type.clone(),
182 snapshot_value,
183 )
184 } else {
185 ColumnDesc::named(col_name.clone(), ColumnId::placeholder(), data_type)
186 };
187
188 column_descs.push(column_desc);
189 if matches!(col.key, ColumnKey::Primary) {
190 pk_names.push(col_name);
191 }
192 }
193
194 if pk_names.is_empty() {
195 return Err(anyhow!("MySQL table doesn't define the primary key").into());
196 }
197 Ok(Self {
198 column_descs,
199 pk_names,
200 })
201 }
202
203 pub fn column_descs(&self) -> &Vec<ColumnDesc> {
204 &self.column_descs
205 }
206
207 pub fn pk_names(&self) -> &Vec<String> {
208 &self.pk_names
209 }
210}
211
212fn derive_default_value(default: ColumnDefault, data_type: &DataType) -> ConnectorResult<Datum> {
213 let datum = match default {
214 ColumnDefault::Null => None,
215 ColumnDefault::Int(val) => match data_type {
216 DataType::Int16 => Some(ScalarImpl::Int16(val as _)),
217 DataType::Int32 => Some(ScalarImpl::Int32(val as _)),
218 DataType::Int64 => Some(ScalarImpl::Int64(val)),
219 DataType::Varchar => {
220 Some(ScalarImpl::from(val.to_string()))
222 }
223 _ => bail!("unexpected default value type for integer"),
224 },
225 ColumnDefault::Real(val) => match data_type {
226 DataType::Float32 => Some(ScalarImpl::Float32(F32::from(val as f32))),
227 DataType::Float64 => Some(ScalarImpl::Float64(val.into())),
228 DataType::Decimal => Some(ScalarImpl::Decimal(
229 Decimal::try_from(val).context("failed to convert default value to decimal")?,
230 )),
231 _ => bail!("unexpected default value type for real"),
232 },
233 ColumnDefault::String(mut val) => {
234 if data_type == &DataType::Timestamptz {
237 val = timestamp_val_to_timestamptz(val.as_str())?;
238 }
239 Some(ScalarImpl::from_text(val.as_str(), data_type).map_err(|e| anyhow!(e)).context(
240 "failed to parse mysql default value expression, only constant is supported",
241 )?)
242 }
243 ColumnDefault::CurrentTimestamp | ColumnDefault::CustomExpr(_) => {
244 bail!("MySQL CURRENT_TIMESTAMP and custom expression default value not supported")
245 }
246 };
247 Ok(datum)
248}
249
250pub fn timestamp_val_to_timestamptz(value_text: &str) -> ConnectorResult<String> {
251 let format = "%Y-%m-%d %H:%M:%S";
252 let naive_datetime = NaiveDateTime::parse_from_str(value_text, format)
253 .map_err(|err| anyhow!("failed to parse mysql timestamp value").context(err))?;
254 let postgres_timestamptz: DateTime<chrono::Utc> =
255 DateTime::<chrono::Utc>::from_naive_utc_and_offset(naive_datetime, chrono::Utc);
256 Ok(postgres_timestamptz
257 .format("%Y-%m-%d %H:%M:%S%:z")
258 .to_string())
259}
260
261pub fn type_name_to_mysql_type(ty_name: &str) -> Option<ColumnType> {
262 macro_rules! column_type {
263 ($($name:literal => $variant:ident),* $(,)?) => {
264 match ty_name.to_lowercase().as_str() {
265 $(
266 $name => Some(ColumnType::$variant(Default::default())),
267 )*
268 "json" => Some(ColumnType::Json),
269 "date" => Some(ColumnType::Date),
270 "bool" => Some(ColumnType::Bool),
271 "tinyblob" => Some(ColumnType::TinyBlob),
272 "mediumblob" => Some(ColumnType::MediumBlob),
273 "longblob" => Some(ColumnType::LongBlob),
274 _ => None,
275 }
276 };
277 }
278
279 column_type! {
280 "bit" => Bit,
281 "tinyint" => TinyInt,
282 "smallint" => SmallInt,
283 "mediumint" => MediumInt,
284 "int" => Int,
285 "bigint" => BigInt,
286 "decimal" => Decimal,
287 "float" => Float,
288 "double" => Double,
289 "time" => Time,
290 "datetime" => DateTime,
291 "timestamp" => Timestamp,
292 "char" => Char,
293 "nchar" => NChar,
294 "varchar" => Varchar,
295 "nvarchar" => NVarchar,
296 "binary" => Binary,
297 "varbinary" => Varbinary,
298 "text" => Text,
299 "tinytext" => TinyText,
300 "mediumtext" => MediumText,
301 "longtext" => LongText,
302 "blob" => Blob,
303 "enum" => Enum,
304 "set" => Set,
305 "geometry" => Geometry,
306 "point" => Point,
307 "linestring" => LineString,
308 "polygon" => Polygon,
309 "multipoint" => MultiPoint,
310 "multilinestring" => MultiLineString,
311 "multipolygon" => MultiPolygon,
312 "geometrycollection" => GeometryCollection,
313 }
314}
315
316pub fn mysql_type_to_rw_type(col_type: &ColumnType) -> ConnectorResult<DataType> {
317 let dtype = match col_type {
318 ColumnType::Serial => DataType::Int32,
319 ColumnType::Bit(attr) => {
320 if let Some(1) = attr.maximum {
321 DataType::Boolean
322 } else {
323 return Err(
324 anyhow!("BIT({}) type not supported", attr.maximum.unwrap_or(0)).into(),
325 );
326 }
327 }
328 ColumnType::TinyInt(_) | ColumnType::SmallInt(_) => DataType::Int16,
329 ColumnType::Bool => DataType::Boolean,
330 ColumnType::MediumInt(_) => DataType::Int32,
331 ColumnType::Int(_) => DataType::Int32,
332 ColumnType::BigInt(_) => DataType::Int64,
333 ColumnType::Decimal(_) => DataType::Decimal,
334 ColumnType::Float(_) => DataType::Float32,
335 ColumnType::Double(_) => DataType::Float64,
336 ColumnType::Date => DataType::Date,
337 ColumnType::Time(_) => DataType::Time,
338 ColumnType::DateTime(_) => DataType::Timestamp,
339 ColumnType::Timestamp(_) => DataType::Timestamptz,
340 ColumnType::Year => DataType::Int32,
341 ColumnType::Char(_)
342 | ColumnType::NChar(_)
343 | ColumnType::Varchar(_)
344 | ColumnType::NVarchar(_) => DataType::Varchar,
345 ColumnType::Binary(_) | ColumnType::Varbinary(_) => DataType::Bytea,
346 ColumnType::Text(_)
347 | ColumnType::TinyText(_)
348 | ColumnType::MediumText(_)
349 | ColumnType::LongText(_) => DataType::Varchar,
350 ColumnType::Blob(_)
351 | ColumnType::TinyBlob
352 | ColumnType::MediumBlob
353 | ColumnType::LongBlob => DataType::Bytea,
354 ColumnType::Enum(_) => DataType::Varchar,
355 ColumnType::Json => DataType::Jsonb,
356 ColumnType::Set(_) => {
357 return Err(anyhow!("SET type not supported").into());
358 }
359 ColumnType::Geometry(_) => {
360 return Err(anyhow!("GEOMETRY type not supported").into());
361 }
362 ColumnType::Point(_) => {
363 return Err(anyhow!("POINT type not supported").into());
364 }
365 ColumnType::LineString(_) => {
366 return Err(anyhow!("LINE string type not supported").into());
367 }
368 ColumnType::Polygon(_) => {
369 return Err(anyhow!("POLYGON type not supported").into());
370 }
371 ColumnType::MultiPoint(_) => {
372 return Err(anyhow!("MULTI POINT type not supported").into());
373 }
374 ColumnType::MultiLineString(_) => {
375 return Err(anyhow!("MULTI LINE STRING type not supported").into());
376 }
377 ColumnType::MultiPolygon(_) => {
378 return Err(anyhow!("MULTI POLYGON type not supported").into());
379 }
380 ColumnType::GeometryCollection(_) => {
381 return Err(anyhow!("GEOMETRY COLLECTION type not supported").into());
382 }
383 ColumnType::Unknown(_) => {
384 return Err(anyhow!("Unknown MySQL data type").into());
385 }
386 };
387
388 Ok(dtype)
389}
390
391pub struct MySqlExternalTableReader {
392 rw_schema: Schema,
393 field_names: String,
394 pool: mysql_async::Pool,
395 upstream_mysql_pk_infos: Vec<(String, String)>, mysql_version: (u8, u8),
397}
398
399impl ExternalTableReader for MySqlExternalTableReader {
400 async fn current_cdc_offset(&self) -> ConnectorResult<CdcOffset> {
401 let mut conn = self.pool.get_conn().await?;
402
403 let sql = if self.is_mysql_8_4_or_later() {
405 "SHOW BINARY LOG STATUS"
406 } else {
407 "SHOW MASTER STATUS"
408 };
409
410 tracing::debug!(
411 "Using SQL command: {} for MySQL version {}.{}",
412 sql,
413 self.mysql_version.0,
414 self.mysql_version.1
415 );
416 let mut rs = conn.query::<mysql_async::Row, _>(sql).await?;
417 let row = rs
418 .iter_mut()
419 .exactly_one()
420 .ok()
421 .context("expect exactly one row when reading binlog offset")?;
422 drop(conn);
423 Ok(CdcOffset::MySql(MySqlOffset {
424 filename: row.take("File").unwrap(),
425 position: row.take("Position").unwrap(),
426 }))
427 }
428
429 fn snapshot_read(
430 &self,
431 table_name: SchemaTableName,
432 start_pk: Option<OwnedRow>,
433 primary_keys: Vec<String>,
434 limit: u32,
435 ) -> BoxStream<'_, ConnectorResult<OwnedRow>> {
436 self.snapshot_read_inner(table_name, start_pk, primary_keys, limit)
437 }
438
439 async fn disconnect(self) -> ConnectorResult<()> {
440 self.pool.disconnect().await.map_err(|e| e.into())
441 }
442
443 fn get_parallel_cdc_splits(
444 &self,
445 _options: CdcTableSnapshotSplitOption,
446 ) -> BoxStream<'_, ConnectorResult<CdcTableSnapshotSplit>> {
447 stream::empty::<ConnectorResult<CdcTableSnapshotSplit>>().boxed()
449 }
450
451 fn split_snapshot_read(
452 &self,
453 _table_name: SchemaTableName,
454 _left: OwnedRow,
455 _right: OwnedRow,
456 _split_columns: Vec<Field>,
457 ) -> BoxStream<'_, ConnectorResult<OwnedRow>> {
458 todo!("implement MySQL CDC parallelized backfill")
459 }
460}
461
462impl MySqlExternalTableReader {
463 async fn get_mysql_version(pool: &mysql_async::Pool) -> ConnectorResult<(u8, u8)> {
465 let mut conn = pool.get_conn().await?;
466 let result: Option<String> = conn.query_first("SELECT VERSION()").await?;
467
468 if let Some(version_str) = result {
469 let parts: Vec<&str> = version_str.split('.').collect();
470 if parts.len() >= 2 {
471 let major_version = parts[0]
472 .parse::<u8>()
473 .context("Failed to parse major version")?;
474 let minor_version = parts[1]
475 .parse::<u8>()
476 .context("Failed to parse minor version")?;
477 return Ok((major_version, minor_version));
478 }
479 }
480 Err(anyhow!("Failed to get MySQL version").into())
481 }
482
483 fn is_mysql_8_4_or_later(&self) -> bool {
485 let (major, minor) = self.mysql_version;
486 major > 8 || (major == 8 && minor >= 4)
487 }
488
489 pub async fn new(config: ExternalTableConfig, rw_schema: Schema) -> ConnectorResult<Self> {
490 let database = config.database.clone();
491 let table = config.table.clone();
492 let pool = build_mysql_connection_pool(
493 &config.host,
494 config.port.parse::<u16>().unwrap(),
495 &config.username,
496 &config.password,
497 &config.database,
498 config.ssl_mode,
499 );
500
501 let field_names = rw_schema
502 .fields
503 .iter()
504 .filter(|f| f.name != CDC_OFFSET_COLUMN_NAME)
505 .map(|f| Self::quote_column(f.name.as_str()))
506 .join(",");
507
508 let upstream_mysql_pk_infos =
510 Self::query_upstream_pk_infos(&pool, &database, &table).await?;
511 let mysql_version = Self::get_mysql_version(&pool).await?;
513 tracing::info!(
514 "MySQL version detected: {}.{}",
515 mysql_version.0,
516 mysql_version.1
517 );
518
519 Ok(Self {
520 rw_schema,
521 field_names,
522 pool,
523 upstream_mysql_pk_infos,
524 mysql_version,
525 })
526 }
527
528 pub fn get_normalized_table_name(table_name: &SchemaTableName) -> String {
529 format!("`{}`.`{}`", table_name.schema_name, table_name.table_name)
531 }
532
533 pub fn get_cdc_offset_parser() -> CdcOffsetParseFunc {
534 Box::new(move |offset| {
535 Ok(CdcOffset::MySql(MySqlOffset::parse_debezium_offset(
536 offset,
537 )?))
538 })
539 }
540
541 async fn query_upstream_pk_infos(
543 pool: &mysql_async::Pool,
544 database: &str,
545 table: &str,
546 ) -> ConnectorResult<Vec<(String, String)>> {
547 let mut conn = pool.get_conn().await?;
548
549 let sql = format!(
551 "SELECT COLUMN_NAME, COLUMN_TYPE
552 FROM INFORMATION_SCHEMA.COLUMNS
553 WHERE TABLE_SCHEMA = '{}'
554 AND TABLE_NAME = '{}'
555 AND COLUMN_KEY = 'PRI'
556 ORDER BY ORDINAL_POSITION",
557 database, table
558 );
559
560 let rs = conn.query::<mysql_async::Row, _>(sql).await?;
561
562 let mut column_infos = Vec::new();
563 for row in &rs {
564 let column_name: String = row.get(0).unwrap();
565 let column_type: String = row.get(1).unwrap();
566 column_infos.push((column_name, column_type));
567 }
568
569 drop(conn);
570
571 Ok(column_infos)
572 }
573
574 fn is_unsigned_type(&self, column_name: &str) -> bool {
576 self.upstream_mysql_pk_infos
577 .iter()
578 .find(|(col_name, _)| col_name == column_name)
579 .map(|(_, col_type)| col_type.to_lowercase().contains("unsigned"))
580 .unwrap_or(false)
581 }
582
583 fn convert_negative_to_unsigned(&self, negative_val: i64) -> u64 {
585 negative_val as u64
586 }
587
588 #[try_stream(boxed, ok = OwnedRow, error = ConnectorError)]
589 async fn snapshot_read_inner(
590 &self,
591 table_name: SchemaTableName,
592 start_pk_row: Option<OwnedRow>,
593 primary_keys: Vec<String>,
594 limit: u32,
595 ) {
596 let order_key = primary_keys
597 .iter()
598 .map(|col| Self::quote_column(col))
599 .join(",");
600 let sql = if start_pk_row.is_none() {
601 format!(
602 "SELECT {} FROM {} ORDER BY {} LIMIT {limit}",
603 self.field_names,
604 Self::get_normalized_table_name(&table_name),
605 order_key,
606 )
607 } else {
608 let filter_expr = Self::filter_expression(&primary_keys);
609 format!(
610 "SELECT {} FROM {} WHERE {} ORDER BY {} LIMIT {limit}",
611 self.field_names,
612 Self::get_normalized_table_name(&table_name),
613 filter_expr,
614 order_key,
615 )
616 };
617 let mut conn = self.pool.get_conn().await?;
618 conn.exec_drop("SET time_zone = \"+00:00\"", ()).await?;
620
621 if let Some(start_pk_row) = start_pk_row {
622 let field_map = self
623 .rw_schema
624 .fields
625 .iter()
626 .map(|f| (f.name.as_str(), f.data_type.clone()))
627 .collect::<HashMap<_, _>>();
628
629 let params: Vec<_> = primary_keys
631 .iter()
632 .zip_eq_fast(start_pk_row.into_iter())
633 .map(|(pk, datum)| {
634 if let Some(value) = datum {
635 let ty = field_map.get(pk.as_str()).unwrap();
636 let val = match ty {
637 DataType::Boolean => Value::from(value.into_bool()),
638 DataType::Int16 => Value::from(value.into_int16()),
639 DataType::Int32 => Value::from(value.into_int32()),
640 DataType::Int64 => {
641 let int64_val = value.into_int64();
642 if int64_val < 0 && self.is_unsigned_type(pk.as_str()) {
643 Value::from(self.convert_negative_to_unsigned(int64_val))
644 } else {
645 Value::from(int64_val)
646 }
647 }
648 DataType::Float32 => Value::from(value.into_float32().into_inner()),
649 DataType::Float64 => Value::from(value.into_float64().into_inner()),
650 DataType::Varchar => Value::from(String::from(value.into_utf8())),
651 DataType::Date => Value::from(value.into_date().0),
652 DataType::Time => Value::from(value.into_time().0),
653 DataType::Timestamp => Value::from(value.into_timestamp().0),
654 DataType::Decimal => Value::from(value.into_decimal().to_string()),
655 DataType::Timestamptz => {
656 let ts = value.into_timestamptz();
659 let datetime_utc = ts.to_datetime_utc();
660 let naive_datetime = datetime_utc.naive_utc();
661 Value::from(naive_datetime)
662 }
663 _ => bail!("unsupported primary key data type: {}", ty),
664 };
665 ConnectorResult::Ok((pk.to_lowercase(), val))
666 } else {
667 bail!("primary key {} cannot be null", pk);
668 }
669 })
670 .try_collect::<_, _, ConnectorError>()?;
671
672 tracing::debug!("snapshot read params: {:?}", ¶ms);
673 let rs_stream = sql
674 .with(Params::from(params))
675 .stream::<mysql_async::Row, _>(&mut conn)
676 .await?;
677
678 let row_stream = rs_stream.map(|row| {
679 let mut row = row?;
681 Ok::<_, ConnectorError>(mysql_row_to_owned_row(&mut row, &self.rw_schema))
682 });
683 pin_mut!(row_stream);
684 #[for_await]
685 for row in row_stream {
686 let row = row?;
687 yield row;
688 }
689 } else {
690 let rs_stream = sql.stream::<mysql_async::Row, _>(&mut conn).await?;
691 let row_stream = rs_stream.map(|row| {
692 let mut row = row?;
694 Ok::<_, ConnectorError>(mysql_row_to_owned_row(&mut row, &self.rw_schema))
695 });
696 pin_mut!(row_stream);
697 #[for_await]
698 for row in row_stream {
699 let row = row?;
700 yield row;
701 }
702 }
703 drop(conn);
704 }
705
706 fn filter_expression(columns: &[String]) -> String {
710 let mut conditions = vec![];
711 conditions.push(format!(
713 "({} > :{})",
714 Self::quote_column(&columns[0]),
715 columns[0].to_lowercase()
716 ));
717 for i in 2..=columns.len() {
718 let mut condition = String::new();
720 for (j, col) in columns.iter().enumerate().take(i - 1) {
721 if j == 0 {
722 condition.push_str(&format!(
723 "({} = :{})",
724 Self::quote_column(col),
725 col.to_lowercase()
726 ));
727 } else {
728 condition.push_str(&format!(
729 " AND ({} = :{})",
730 Self::quote_column(col),
731 col.to_lowercase()
732 ));
733 }
734 }
735 condition.push_str(&format!(
737 " AND ({} > :{})",
738 Self::quote_column(&columns[i - 1]),
739 columns[i - 1].to_lowercase()
740 ));
741 conditions.push(format!("({})", condition));
742 }
743 if columns.len() > 1 {
744 conditions.join(" OR ")
745 } else {
746 conditions.join("")
747 }
748 }
749
750 fn quote_column(column: &str) -> String {
751 format!("`{}`", column)
752 }
753}
754
755#[cfg(test)]
756mod tests {
757 use std::collections::HashMap;
758
759 use futures::pin_mut;
760 use futures_async_stream::for_await;
761 use maplit::{convert_args, hashmap};
762 use risingwave_common::catalog::{ColumnDesc, ColumnId, Field, Schema};
763 use risingwave_common::types::DataType;
764
765 use crate::source::cdc::external::mysql::MySqlExternalTable;
766 use crate::source::cdc::external::{
767 CdcOffset, ExternalTableConfig, ExternalTableReader, MySqlExternalTableReader, MySqlOffset,
768 SchemaTableName,
769 };
770
771 #[ignore]
772 #[tokio::test]
773 async fn test_mysql_schema() {
774 let config = ExternalTableConfig {
775 connector: "mysql-cdc".to_owned(),
776 host: "localhost".to_owned(),
777 port: "8306".to_owned(),
778 username: "root".to_owned(),
779 password: "123456".to_owned(),
780 database: "mydb".to_owned(),
781 schema: "".to_owned(),
782 table: "part".to_owned(),
783 ssl_mode: Default::default(),
784 ssl_root_cert: None,
785 encrypt: "false".to_owned(),
786 };
787
788 let table = MySqlExternalTable::connect(config).await.unwrap();
789 println!("columns: {:?}", &table.column_descs);
790 println!("primary keys: {:?}", &table.pk_names);
791 }
792
793 #[test]
794 fn test_mysql_filter_expr() {
795 let cols = vec!["id".to_owned()];
796 let expr = MySqlExternalTableReader::filter_expression(&cols);
797 assert_eq!(expr, "(`id` > :id)");
798
799 let cols = vec!["aa".to_owned(), "bb".to_owned(), "cc".to_owned()];
800 let expr = MySqlExternalTableReader::filter_expression(&cols);
801 assert_eq!(
802 expr,
803 "(`aa` > :aa) OR ((`aa` = :aa) AND (`bb` > :bb)) OR ((`aa` = :aa) AND (`bb` = :bb) AND (`cc` > :cc))"
804 );
805 }
806
807 #[test]
808 fn test_mysql_binlog_offset() {
809 let off0_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000001", "pos": 105622, "snapshot": true }, "isHeartbeat": false }"#;
810 let off1_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000007", "pos": 1062363217, "snapshot": true }, "isHeartbeat": false }"#;
811 let off2_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000007", "pos": 659687560, "snapshot": true }, "isHeartbeat": false }"#;
812 let off3_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000008", "pos": 7665875, "snapshot": true }, "isHeartbeat": false }"#;
813 let off4_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000008", "pos": 7665875, "snapshot": true }, "isHeartbeat": false }"#;
814
815 let off0 = CdcOffset::MySql(MySqlOffset::parse_debezium_offset(off0_str).unwrap());
816 let off1 = CdcOffset::MySql(MySqlOffset::parse_debezium_offset(off1_str).unwrap());
817 let off2 = CdcOffset::MySql(MySqlOffset::parse_debezium_offset(off2_str).unwrap());
818 let off3 = CdcOffset::MySql(MySqlOffset::parse_debezium_offset(off3_str).unwrap());
819 let off4 = CdcOffset::MySql(MySqlOffset::parse_debezium_offset(off4_str).unwrap());
820
821 assert!(off0 <= off1);
822 assert!(off1 > off2);
823 assert!(off2 < off3);
824 assert_eq!(off3, off4);
825 }
826
827 #[ignore]
829 #[tokio::test]
830 async fn test_mysql_table_reader() {
831 let columns = [
832 ColumnDesc::named("v1", ColumnId::new(1), DataType::Int32),
833 ColumnDesc::named("v2", ColumnId::new(2), DataType::Decimal),
834 ColumnDesc::named("v3", ColumnId::new(3), DataType::Varchar),
835 ColumnDesc::named("v4", ColumnId::new(4), DataType::Date),
836 ];
837 let rw_schema = Schema {
838 fields: columns.iter().map(Field::from).collect(),
839 };
840 let props: HashMap<String, String> = convert_args!(hashmap!(
841 "hostname" => "localhost",
842 "port" => "8306",
843 "username" => "root",
844 "password" => "123456",
845 "database.name" => "mytest",
846 "table.name" => "t1"));
847
848 let config =
849 serde_json::from_value::<ExternalTableConfig>(serde_json::to_value(props).unwrap())
850 .unwrap();
851 let reader = MySqlExternalTableReader::new(config, rw_schema)
852 .await
853 .unwrap();
854 let offset = reader.current_cdc_offset().await.unwrap();
855 println!("BinlogOffset: {:?}", offset);
856
857 let off0_str = r#"{ "sourcePartition": { "server": "test" }, "sourceOffset": { "ts_sec": 1670876905, "file": "binlog.000001", "pos": 105622, "snapshot": true }, "isHeartbeat": false }"#;
858 let parser = MySqlExternalTableReader::get_cdc_offset_parser();
859 println!("parsed offset: {:?}", parser(off0_str).unwrap());
860 let table_name = SchemaTableName {
861 schema_name: "mytest".to_owned(),
862 table_name: "t1".to_owned(),
863 };
864
865 let stream = reader.snapshot_read(table_name, None, vec!["v1".to_owned()], 1000);
866 pin_mut!(stream);
867 #[for_await]
868 for row in stream {
869 println!("OwnedRow: {:?}", row);
870 }
871 }
872}