1pub mod mock_external_table;
16pub mod postgres;
17pub mod sql_server;
18
19pub mod mysql;
20
21use std::collections::{BTreeMap, HashMap};
22
23use anyhow::{Context, anyhow};
24use futures::pin_mut;
25use futures::stream::BoxStream;
26use futures_async_stream::try_stream;
27use risingwave_common::bail;
28use risingwave_common::catalog::{ColumnDesc, Field, Schema};
29use risingwave_common::row::OwnedRow;
30use risingwave_common::secret::LocalSecretManager;
31use risingwave_pb::catalog::table::CdcTableType as PbCdcTableType;
32use risingwave_pb::secret::PbSecretRef;
33use serde::{Deserialize, Serialize};
34
35use crate::WithPropertiesExt;
36use crate::connector_common::{PgConnectionConfig, PostgresExternalTable, SslMode};
37use crate::error::{ConnectorError, ConnectorResult};
38use crate::parser::mysql_row_to_owned_row;
39use crate::source::CdcTableSnapshotSplit;
40use crate::source::cdc::CdcSourceType;
41use crate::source::cdc::external::mock_external_table::MockExternalTableReader;
42use crate::source::cdc::external::mysql::{
43 MySqlExternalTable, MySqlExternalTableReader, MySqlOffset,
44};
45use crate::source::cdc::external::postgres::{PostgresExternalTableReader, PostgresOffset};
46use crate::source::cdc::external::sql_server::{
47 SqlServerExternalTable, SqlServerExternalTableReader, SqlServerOffset,
48};
49
50#[derive(Debug, Clone, PartialEq, Eq, Hash)]
51pub enum ExternalCdcTableType {
52 Undefined,
53 Mock,
54 MySql,
55 Postgres,
56 SqlServer,
57 Citus,
58 Mongo,
59}
60
61impl ExternalCdcTableType {
62 pub fn from_properties(with_properties: &impl WithPropertiesExt) -> Self {
63 let connector = with_properties.get_connector().unwrap_or_default();
64 match connector.as_str() {
65 "mysql-cdc" => Self::MySql,
66 "postgres-cdc" => Self::Postgres,
67 "citus-cdc" => Self::Citus,
68 "sqlserver-cdc" => Self::SqlServer,
69 "mongodb-cdc" => Self::Mongo,
70 _ => Self::Undefined,
71 }
72 }
73
74 pub fn can_backfill(&self) -> bool {
75 matches!(self, Self::MySql | Self::Postgres | Self::SqlServer)
76 }
77
78 pub fn enable_transaction_metadata(&self) -> bool {
79 matches!(self, Self::MySql | Self::Postgres)
83 }
84
85 pub async fn create_table_reader(
86 &self,
87 config: ExternalTableConfig,
88 schema: Schema,
89 pk_indices: Vec<usize>,
90 schema_table_name: SchemaTableName,
91 ) -> ConnectorResult<ExternalTableReaderImpl> {
92 match self {
93 Self::MySql => Ok(ExternalTableReaderImpl::MySql(
94 MySqlExternalTableReader::new(config, schema).await?,
95 )),
96 Self::Postgres => Ok(ExternalTableReaderImpl::Postgres(
97 PostgresExternalTableReader::new(config, schema, pk_indices, schema_table_name)
98 .await?,
99 )),
100 Self::SqlServer => Ok(ExternalTableReaderImpl::SqlServer(
101 SqlServerExternalTableReader::new(config, schema, pk_indices).await?,
102 )),
103 Self::Mock => Ok(ExternalTableReaderImpl::Mock(MockExternalTableReader::new())),
105 _ => bail!("invalid external table type: {:?}", *self),
106 }
107 }
108}
109
110impl From<ExternalCdcTableType> for PbCdcTableType {
111 fn from(cdc_table_type: ExternalCdcTableType) -> Self {
112 match cdc_table_type {
113 ExternalCdcTableType::Postgres => Self::Postgres,
114 ExternalCdcTableType::MySql => Self::Mysql,
115 ExternalCdcTableType::SqlServer => Self::Sqlserver,
116
117 ExternalCdcTableType::Citus => Self::Citus,
118 ExternalCdcTableType::Mongo => Self::Mongo,
119 ExternalCdcTableType::Undefined | ExternalCdcTableType::Mock => Self::Unspecified,
120 }
121 }
122}
123
124impl From<PbCdcTableType> for ExternalCdcTableType {
125 fn from(cdc_table_type: PbCdcTableType) -> Self {
126 match cdc_table_type {
127 PbCdcTableType::Postgres => Self::Postgres,
128 PbCdcTableType::Mysql => Self::MySql,
129 PbCdcTableType::Sqlserver => Self::SqlServer,
130 PbCdcTableType::Mongo => Self::Mongo,
131 PbCdcTableType::Citus => Self::Citus,
132 PbCdcTableType::Unspecified => Self::Undefined,
133 }
134 }
135}
136
137#[derive(Debug, Clone, PartialEq)]
138pub struct SchemaTableName {
139 pub schema_name: String,
141 pub table_name: String,
142}
143
144pub const TABLE_NAME_KEY: &str = "table.name";
145pub const SCHEMA_NAME_KEY: &str = "schema.name";
146pub const DATABASE_NAME_KEY: &str = "database.name";
147
148impl SchemaTableName {
149 pub fn from_properties(properties: &BTreeMap<String, String>) -> Self {
150 let table_type = ExternalCdcTableType::from_properties(properties);
151 let table_name = properties.get(TABLE_NAME_KEY).cloned().unwrap_or_default();
152
153 let schema_name = match table_type {
154 ExternalCdcTableType::MySql => properties
155 .get(DATABASE_NAME_KEY)
156 .cloned()
157 .unwrap_or_default(),
158 ExternalCdcTableType::Postgres | ExternalCdcTableType::Citus => {
159 properties.get(SCHEMA_NAME_KEY).cloned().unwrap_or_default()
160 }
161 ExternalCdcTableType::SqlServer => {
162 properties.get(SCHEMA_NAME_KEY).cloned().unwrap_or_default()
163 }
164 _ => {
165 unreachable!("invalid external table type: {:?}", table_type);
166 }
167 };
168
169 Self {
170 schema_name,
171 table_name,
172 }
173 }
174}
175
176#[derive(Debug, Clone, PartialEq, PartialOrd, Serialize, Deserialize)]
177pub enum CdcOffset {
178 MySql(MySqlOffset),
179 Postgres(PostgresOffset),
180 SqlServer(SqlServerOffset),
181}
182
183#[derive(Debug, Clone, Serialize, Deserialize)]
199pub struct DebeziumOffset {
200 #[serde(rename = "sourcePartition")]
201 pub source_partition: HashMap<String, String>,
202 #[serde(rename = "sourceOffset")]
203 pub source_offset: DebeziumSourceOffset,
204 #[serde(rename = "isHeartbeat")]
205 pub is_heartbeat: bool,
206}
207
208#[derive(Debug, Default, Clone, Serialize, Deserialize)]
209pub struct DebeziumSourceOffset {
210 pub last_snapshot_record: Option<bool>,
212 pub snapshot: Option<bool>,
214
215 pub file: Option<String>,
217 pub pos: Option<u64>,
218
219 pub lsn: Option<u64>,
221 #[serde(rename = "txId")]
222 pub txid: Option<i64>,
223 pub tx_usec: Option<u64>,
224 pub lsn_commit: Option<u64>,
225 pub lsn_proc: Option<u64>,
226
227 pub commit_lsn: Option<String>,
229 pub change_lsn: Option<String>,
230}
231
232pub type CdcOffsetParseFunc = Box<dyn Fn(&str) -> ConnectorResult<CdcOffset> + Send>;
233
234pub trait ExternalTableReader: Sized {
235 async fn current_cdc_offset(&self) -> ConnectorResult<CdcOffset>;
236
237 async fn disconnect(self) -> ConnectorResult<()> {
240 Ok(())
241 }
242
243 fn snapshot_read(
244 &self,
245 table_name: SchemaTableName,
246 start_pk: Option<OwnedRow>,
247 primary_keys: Vec<String>,
248 limit: u32,
249 ) -> BoxStream<'_, ConnectorResult<OwnedRow>>;
250
251 fn get_parallel_cdc_splits(
252 &self,
253 options: CdcTableSnapshotSplitOption,
254 ) -> BoxStream<'_, ConnectorResult<CdcTableSnapshotSplit>>;
255
256 fn split_snapshot_read(
257 &self,
258 table_name: SchemaTableName,
259 left: OwnedRow,
260 right: OwnedRow,
261 split_columns: Vec<Field>,
262 ) -> BoxStream<'_, ConnectorResult<OwnedRow>>;
263}
264
265pub struct CdcTableSnapshotSplitOption {
266 pub backfill_num_rows_per_split: u64,
267 pub backfill_as_even_splits: bool,
268 pub backfill_split_pk_column_index: u32,
269}
270
271pub enum ExternalTableReaderImpl {
272 MySql(MySqlExternalTableReader),
273 Postgres(PostgresExternalTableReader),
274 SqlServer(SqlServerExternalTableReader),
275 Mock(MockExternalTableReader),
276}
277
278#[derive(Debug, Default, Clone, Deserialize)]
279pub struct ExternalTableConfig {
280 pub connector: String,
281
282 #[serde(rename = "hostname")]
283 pub host: String,
284 pub port: String,
285 pub username: String,
286 pub password: String,
287 #[serde(rename = "database.name")]
288 pub database: String,
289 #[serde(rename = "schema.name", default = "Default::default")]
290 pub schema: String,
291 #[serde(rename = "table.name")]
292 pub table: String,
293 #[serde(rename = "ssl.mode", default = "postgres_ssl_mode_default")]
297 #[serde(alias = "debezium.database.sslmode")]
298 pub ssl_mode: SslMode,
299
300 #[serde(rename = "ssl.root.cert")]
301 #[serde(alias = "debezium.database.sslrootcert")]
302 pub ssl_root_cert: Option<String>,
303
304 #[serde(rename = "database.encrypt", default = "Default::default")]
307 pub encrypt: String,
308}
309
310fn postgres_ssl_mode_default() -> SslMode {
311 SslMode::Disabled
313}
314
315impl ExternalTableConfig {
316 pub fn try_from_btreemap(
317 connect_properties: BTreeMap<String, String>,
318 secret_refs: BTreeMap<String, PbSecretRef>,
319 ) -> ConnectorResult<Self> {
320 let options_with_secret =
321 LocalSecretManager::global().fill_secrets(connect_properties, secret_refs)?;
322 let json_value = serde_json::to_value(options_with_secret)?;
323 let config = serde_json::from_value::<ExternalTableConfig>(json_value)?;
324 Ok(config)
325 }
326
327 pub fn pg_connection_config(&self) -> ConnectorResult<PgConnectionConfig> {
331 let port = self
332 .port
333 .parse::<u16>()
334 .with_context(|| format!("invalid postgres port `{}`", self.port))?;
335 Ok(PgConnectionConfig {
336 host: self.host.clone(),
337 port,
338 user: self.username.clone(),
339 password: self.password.clone(),
340 database: self.database.clone(),
341 ssl_mode: self.ssl_mode.clone(),
342 ssl_root_cert: self.ssl_root_cert.clone(),
343 })
344 }
345}
346
347impl ExternalTableReader for ExternalTableReaderImpl {
348 async fn current_cdc_offset(&self) -> ConnectorResult<CdcOffset> {
349 match self {
350 ExternalTableReaderImpl::MySql(mysql) => mysql.current_cdc_offset().await,
351 ExternalTableReaderImpl::Postgres(postgres) => postgres.current_cdc_offset().await,
352 ExternalTableReaderImpl::SqlServer(sql_server) => sql_server.current_cdc_offset().await,
353 ExternalTableReaderImpl::Mock(mock) => mock.current_cdc_offset().await,
354 }
355 }
356
357 fn snapshot_read(
358 &self,
359 table_name: SchemaTableName,
360 start_pk: Option<OwnedRow>,
361 primary_keys: Vec<String>,
362 limit: u32,
363 ) -> BoxStream<'_, ConnectorResult<OwnedRow>> {
364 self.snapshot_read_inner(table_name, start_pk, primary_keys, limit)
365 }
366
367 fn get_parallel_cdc_splits(
368 &self,
369 options: CdcTableSnapshotSplitOption,
370 ) -> BoxStream<'_, ConnectorResult<CdcTableSnapshotSplit>> {
371 self.get_parallel_cdc_splits_inner(options)
372 }
373
374 fn split_snapshot_read(
375 &self,
376 table_name: SchemaTableName,
377 left: OwnedRow,
378 right: OwnedRow,
379 split_columns: Vec<Field>,
380 ) -> BoxStream<'_, ConnectorResult<OwnedRow>> {
381 self.split_snapshot_read_inner(table_name, left, right, split_columns)
382 }
383}
384
385impl ExternalTableReaderImpl {
386 pub fn pk_column_unsigned_i64_compare_flags(
390 &self,
391 pk_names: &[String],
392 ) -> ConnectorResult<Vec<bool>> {
393 match self {
394 ExternalTableReaderImpl::MySql(mysql) => {
395 mysql.pk_column_unsigned_i64_compare_flags(pk_names)
396 }
397 _ => Ok(vec![false; pk_names.len()]),
398 }
399 }
400
401 pub fn get_cdc_offset_parser(&self) -> CdcOffsetParseFunc {
402 match self {
403 ExternalTableReaderImpl::MySql(_) => MySqlExternalTableReader::get_cdc_offset_parser(),
404 ExternalTableReaderImpl::Postgres(_) => {
405 PostgresExternalTableReader::get_cdc_offset_parser()
406 }
407 ExternalTableReaderImpl::SqlServer(_) => {
408 SqlServerExternalTableReader::get_cdc_offset_parser()
409 }
410 ExternalTableReaderImpl::Mock(_) => MockExternalTableReader::get_cdc_offset_parser(),
411 }
412 }
413
414 #[try_stream(boxed, ok = OwnedRow, error = ConnectorError)]
415 async fn snapshot_read_inner(
416 &self,
417 table_name: SchemaTableName,
418 start_pk: Option<OwnedRow>,
419 primary_keys: Vec<String>,
420 limit: u32,
421 ) {
422 let stream = match self {
423 ExternalTableReaderImpl::MySql(mysql) => {
424 mysql.snapshot_read(table_name, start_pk, primary_keys, limit)
425 }
426 ExternalTableReaderImpl::Postgres(postgres) => {
427 postgres.snapshot_read(table_name, start_pk, primary_keys, limit)
428 }
429 ExternalTableReaderImpl::SqlServer(sql_server) => {
430 sql_server.snapshot_read(table_name, start_pk, primary_keys, limit)
431 }
432 ExternalTableReaderImpl::Mock(mock) => {
433 mock.snapshot_read(table_name, start_pk, primary_keys, limit)
434 }
435 };
436
437 pin_mut!(stream);
438 #[for_await]
439 for row in stream {
440 let row = row?;
441 yield row;
442 }
443 }
444
445 #[try_stream(boxed, ok = CdcTableSnapshotSplit, error = ConnectorError)]
446 async fn get_parallel_cdc_splits_inner(&self, options: CdcTableSnapshotSplitOption) {
447 let stream = match self {
448 ExternalTableReaderImpl::MySql(e) => e.get_parallel_cdc_splits(options),
449 ExternalTableReaderImpl::Postgres(e) => e.get_parallel_cdc_splits(options),
450 ExternalTableReaderImpl::SqlServer(e) => e.get_parallel_cdc_splits(options),
451 ExternalTableReaderImpl::Mock(e) => e.get_parallel_cdc_splits(options),
452 };
453 pin_mut!(stream);
454 #[for_await]
455 for row in stream {
456 let row = row?;
457 yield row;
458 }
459 }
460
461 #[try_stream(boxed, ok = OwnedRow, error = ConnectorError)]
462 async fn split_snapshot_read_inner(
463 &self,
464 table_name: SchemaTableName,
465 left: OwnedRow,
466 right: OwnedRow,
467 split_columns: Vec<Field>,
468 ) {
469 let stream = match self {
470 ExternalTableReaderImpl::MySql(mysql) => {
471 mysql.split_snapshot_read(table_name, left, right, split_columns)
472 }
473 ExternalTableReaderImpl::Postgres(postgres) => {
474 postgres.split_snapshot_read(table_name, left, right, split_columns)
475 }
476 ExternalTableReaderImpl::SqlServer(sql_server) => {
477 sql_server.split_snapshot_read(table_name, left, right, split_columns)
478 }
479 ExternalTableReaderImpl::Mock(mock) => {
480 mock.split_snapshot_read(table_name, left, right, split_columns)
481 }
482 };
483
484 pin_mut!(stream);
485 #[for_await]
486 for row in stream {
487 let row = row?;
488 yield row;
489 }
490 }
491}
492
493pub enum ExternalTableImpl {
494 MySql(MySqlExternalTable),
495 Postgres(PostgresExternalTable),
496 SqlServer(SqlServerExternalTable),
497}
498
499impl ExternalTableImpl {
500 pub async fn connect(config: ExternalTableConfig) -> ConnectorResult<Self> {
501 let cdc_source_type = CdcSourceType::from(config.connector.as_str());
502 match cdc_source_type {
503 CdcSourceType::Mysql => Ok(ExternalTableImpl::MySql(
504 MySqlExternalTable::connect(config).await?,
505 )),
506 CdcSourceType::Postgres => {
507 let pg_conn = config.pg_connection_config()?;
508 Ok(ExternalTableImpl::Postgres(
509 PostgresExternalTable::connect(&pg_conn, &config.schema, &config.table, false)
510 .await?,
511 ))
512 }
513 CdcSourceType::SqlServer => Ok(ExternalTableImpl::SqlServer(
514 SqlServerExternalTable::connect(config).await?,
515 )),
516 _ => Err(anyhow!("Unsupported cdc connector type: {}", config.connector).into()),
517 }
518 }
519
520 pub fn column_descs(&self) -> &Vec<ColumnDesc> {
521 match self {
522 ExternalTableImpl::MySql(mysql) => mysql.column_descs(),
523 ExternalTableImpl::Postgres(postgres) => postgres.column_descs(),
524 ExternalTableImpl::SqlServer(sql_server) => sql_server.column_descs(),
525 }
526 }
527
528 pub fn pk_names(&self) -> &Vec<String> {
529 match self {
530 ExternalTableImpl::MySql(mysql) => mysql.pk_names(),
531 ExternalTableImpl::Postgres(postgres) => postgres.pk_names(),
532 ExternalTableImpl::SqlServer(sql_server) => sql_server.pk_names(),
533 }
534 }
535}
536
537pub const CDC_TABLE_SPLIT_ID_START: i64 = 1;