Skip to main content

risingwave_connector/parser/debezium/
simd_json_parser.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::fmt::Debug;
16
17use anyhow::Context;
18use simd_json::BorrowedValue;
19use simd_json::prelude::MutableObject;
20
21use crate::error::ConnectorResult;
22use crate::parser::unified::AccessImpl;
23use crate::parser::unified::debezium::MongoJsonAccess;
24use crate::parser::unified::json::{
25    BigintUnsignedHandlingMode, JsonAccess, JsonParseOptions, TimeHandling, TimestampHandling,
26    TimestamptzHandling,
27};
28use crate::parser::{AccessBuilder, MongoProperties};
29
30#[derive(Debug)]
31pub struct DebeziumJsonAccessBuilder {
32    value: Option<Vec<u8>>,
33    json_parse_options: JsonParseOptions,
34}
35
36impl DebeziumJsonAccessBuilder {
37    pub fn new(
38        timestamptz_handling: TimestamptzHandling,
39        timestamp_handling: TimestampHandling,
40        time_handling: TimeHandling,
41        bigint_unsigned_handling: BigintUnsignedHandlingMode,
42        handle_toast_columns: bool,
43    ) -> ConnectorResult<Self> {
44        Ok(Self {
45            value: None,
46            json_parse_options: JsonParseOptions::new_for_debezium(
47                timestamptz_handling,
48                timestamp_handling,
49                time_handling,
50                bigint_unsigned_handling,
51                handle_toast_columns,
52            ),
53        })
54    }
55
56    pub fn new_for_schema_event() -> ConnectorResult<Self> {
57        Ok(Self {
58            value: None,
59            json_parse_options: JsonParseOptions::default(),
60        })
61    }
62}
63
64impl AccessBuilder for DebeziumJsonAccessBuilder {
65    async fn generate_accessor(
66        &mut self,
67        payload: Vec<u8>,
68        _: &crate::source::SourceMeta,
69    ) -> ConnectorResult<AccessImpl<'_>> {
70        self.value = Some(payload);
71        let mut event: BorrowedValue<'_> =
72            simd_json::to_borrowed_value(self.value.as_mut().unwrap())
73                .context("failed to parse debezium json payload")?;
74
75        let payload = if let Some(payload) = event.get_mut("payload") {
76            std::mem::take(payload)
77        } else {
78            event
79        };
80
81        Ok(AccessImpl::Json(JsonAccess::new_with_options(
82            payload,
83            &self.json_parse_options,
84        )))
85    }
86}
87
88#[derive(Debug)]
89pub struct DebeziumMongoJsonAccessBuilder {
90    value: Option<Vec<u8>>,
91    json_parse_options: JsonParseOptions,
92    strong_schema: bool,
93}
94
95impl DebeziumMongoJsonAccessBuilder {
96    pub fn new(props: MongoProperties) -> anyhow::Result<Self> {
97        Ok(Self {
98            value: None,
99            json_parse_options: JsonParseOptions::new_for_debezium(
100                TimestamptzHandling::GuessNumberUnit,
101                TimestampHandling::GuessNumberUnit,
102                TimeHandling::Micro,
103                BigintUnsignedHandlingMode::Long,
104                false,
105            ),
106            strong_schema: props.strong_schema,
107        })
108    }
109}
110
111impl AccessBuilder for DebeziumMongoJsonAccessBuilder {
112    async fn generate_accessor(
113        &mut self,
114        payload: Vec<u8>,
115        _: &crate::source::SourceMeta,
116    ) -> ConnectorResult<AccessImpl<'_>> {
117        self.value = Some(payload);
118        let mut event: BorrowedValue<'_> =
119            simd_json::to_borrowed_value(self.value.as_mut().unwrap())
120                .context("failed to parse debezium mongo json payload")?;
121
122        let payload = if let Some(payload) = event.get_mut("payload") {
123            std::mem::take(payload)
124        } else {
125            event
126        };
127
128        Ok(AccessImpl::MongoJson(MongoJsonAccess::new(
129            JsonAccess::new_with_options(payload, &self.json_parse_options),
130            self.strong_schema,
131        )))
132    }
133}
134
135#[cfg(test)]
136mod tests {
137    use chrono::{NaiveDate, NaiveTime};
138    use risingwave_common::array::{Op, StructValue};
139    use risingwave_common::catalog::ColumnId;
140    use risingwave_common::row::{OwnedRow, Row};
141    use risingwave_common::types::{DataType, Date, Interval, Scalar, ScalarImpl, Time, Timestamp};
142    use serde_json::Value;
143    use thiserror_ext::AsReport;
144
145    use crate::parser::{
146        DebeziumParser, DebeziumProps, EncodingProperties, JsonProperties, ProtocolProperties,
147        SourceColumnDesc, SourceStreamChunkBuilder, SpecificParserConfig,
148    };
149    use crate::source::{SourceContext, SourceCtrlOpts};
150
151    fn assert_json_eq(parse_result: &Option<ScalarImpl>, json_str: &str) {
152        if let Some(ScalarImpl::Jsonb(json_val)) = parse_result {
153            let mut json_string = String::new();
154            json_val
155                .as_scalar_ref()
156                .force_str(&mut json_string)
157                .unwrap();
158            let val1: Value = serde_json::from_str(json_string.as_str()).unwrap();
159            let val2: Value = serde_json::from_str(json_str).unwrap();
160            assert_eq!(val1, val2);
161        }
162    }
163
164    async fn build_parser(rw_columns: Vec<SourceColumnDesc>) -> DebeziumParser {
165        let props = SpecificParserConfig {
166            encoding_config: EncodingProperties::Json(JsonProperties {
167                use_schema_registry: false,
168                timestamptz_handling: None,
169                timestamp_handling: None,
170                time_handling: None,
171                bigint_unsigned_handling: None,
172                handle_toast_columns: false,
173            }),
174            protocol_config: ProtocolProperties::Debezium(DebeziumProps::default()),
175        };
176        DebeziumParser::new(props, rw_columns, SourceContext::dummy().into())
177            .await
178            .unwrap()
179    }
180
181    async fn parse_one(
182        mut parser: DebeziumParser,
183        columns: Vec<SourceColumnDesc>,
184        payload: Vec<u8>,
185    ) -> Vec<(Op, OwnedRow)> {
186        let mut builder = SourceStreamChunkBuilder::new(columns, SourceCtrlOpts::for_test());
187        parser
188            .parse_inner(None, Some(payload), builder.row_writer())
189            .await
190            .unwrap();
191        builder.finish_current_chunk();
192        let chunk = builder.consume_ready_chunks().next().unwrap();
193        chunk
194            .rows()
195            .map(|(op, row_ref)| (op, row_ref.into_owned_row()))
196            .collect::<Vec<_>>()
197    }
198
199    mod test1_basic {
200        use super::*;
201
202        fn get_test1_columns() -> Vec<SourceColumnDesc> {
203            vec![
204                SourceColumnDesc::simple("id", DataType::Int32, ColumnId::from(0)),
205                SourceColumnDesc::simple("name", DataType::Varchar, ColumnId::from(1)),
206                SourceColumnDesc::simple("description", DataType::Varchar, ColumnId::from(2)),
207                SourceColumnDesc::simple("weight", DataType::Float64, ColumnId::from(3)),
208            ]
209        }
210
211        #[tokio::test]
212        async fn test1_debezium_json_parser_read() {
213            //     "before": null,
214            //     "after": {
215            //       "id": 101,
216            //       "name": "scooter",
217            //       "description": "Small 2-wheel scooter",
218            //       "weight": 1.234
219            //     },
220            let input = vec![
221                // data with payload field
222                br#"{"payload":{"before":null,"after":{"id":101,"name":"scooter","description":"Small 2-wheel scooter","weight":1.234},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639547113601,"snapshot":"true","db":"inventory","sequence":null,"table":"products","server_id":0,"gtid":null,"file":"mysql-bin.000003","pos":156,"row":0,"thread":null,"query":null},"op":"r","ts_ms":1639547113602,"transaction":null}}"#.to_vec(),
223                // data without payload field
224                br#"{"before":null,"after":{"id":101,"name":"scooter","description":"Small 2-wheel scooter","weight":1.234},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639547113601,"snapshot":"true","db":"inventory","sequence":null,"table":"products","server_id":0,"gtid":null,"file":"mysql-bin.000003","pos":156,"row":0,"thread":null,"query":null},"op":"r","ts_ms":1639547113602,"transaction":null}"#.to_vec()];
225
226            let columns = get_test1_columns();
227
228            for data in input {
229                let parser = build_parser(columns.clone()).await;
230                let [(_op, row)]: [_; 1] = parse_one(parser, columns.clone(), data)
231                    .await
232                    .try_into()
233                    .unwrap();
234
235                assert!(row[0].eq(&Some(ScalarImpl::Int32(101))));
236                assert!(row[1].eq(&Some(ScalarImpl::Utf8("scooter".into()))));
237                assert!(row[2].eq(&Some(ScalarImpl::Utf8("Small 2-wheel scooter".into()))));
238                assert!(row[3].eq(&Some(ScalarImpl::Float64(1.234.into()))));
239            }
240        }
241
242        #[tokio::test]
243        async fn test1_debezium_json_parser_insert() {
244            //     "before": null,
245            //     "after": {
246            //       "id": 102,
247            //       "name": "car battery",
248            //       "description": "12V car battery",
249            //       "weight": 8.1
250            //     },
251            let input = vec![
252                // data with payload field
253                br#"{"payload":{"before":null,"after":{"id":102,"name":"car battery","description":"12V car battery","weight":8.1},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551564000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":717,"row":0,"thread":null,"query":null},"op":"c","ts_ms":1639551564960,"transaction":null}}"#.to_vec(),
254                // data without payload field
255                br#"{"before":null,"after":{"id":102,"name":"car battery","description":"12V car battery","weight":8.1},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551564000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":717,"row":0,"thread":null,"query":null},"op":"c","ts_ms":1639551564960,"transaction":null}"#.to_vec()];
256
257            let columns = get_test1_columns();
258
259            for data in input {
260                let parser = build_parser(columns.clone()).await;
261                let [(op, row)]: [_; 1] = parse_one(parser, columns.clone(), data)
262                    .await
263                    .try_into()
264                    .unwrap();
265                assert_eq!(op, Op::Insert);
266
267                assert!(row[0].eq(&Some(ScalarImpl::Int32(102))));
268                assert!(row[1].eq(&Some(ScalarImpl::Utf8("car battery".into()))));
269                assert!(row[2].eq(&Some(ScalarImpl::Utf8("12V car battery".into()))));
270                assert!(row[3].eq(&Some(ScalarImpl::Float64(8.1.into()))));
271            }
272        }
273
274        #[tokio::test]
275        async fn test1_debezium_json_parser_delete() {
276            //     "before": {
277            //       "id": 101,
278            //       "name": "scooter",
279            //       "description": "Small 2-wheel scooter",
280            //       "weight": 1.234
281            //     },
282            //     "after": null,
283            let input = vec![
284                // data with payload field
285                br#"{"payload":{"before":{"id":101,"name":"scooter","description":"Small 2-wheel scooter","weight":1.234},"after":null,"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551767000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1045,"row":0,"thread":null,"query":null},"op":"d","ts_ms":1639551767775,"transaction":null}}"#.to_vec(),
286                // data without payload field
287                br#"{"before":{"id":101,"name":"scooter","description":"Small 2-wheel scooter","weight":1.234},"after":null,"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551767000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1045,"row":0,"thread":null,"query":null},"op":"d","ts_ms":1639551767775,"transaction":null}"#.to_vec()];
288
289            for data in input {
290                let columns = get_test1_columns();
291                let parser = build_parser(columns.clone()).await;
292                let [(op, row)]: [_; 1] = parse_one(parser, columns.clone(), data)
293                    .await
294                    .try_into()
295                    .unwrap();
296
297                assert_eq!(op, Op::Delete);
298
299                assert!(row[0].eq(&Some(ScalarImpl::Int32(101))));
300                assert!(row[1].eq(&Some(ScalarImpl::Utf8("scooter".into()))));
301                assert!(row[2].eq(&Some(ScalarImpl::Utf8("Small 2-wheel scooter".into()))));
302                assert!(row[3].eq(&Some(ScalarImpl::Float64(1.234.into()))));
303            }
304        }
305
306        #[tokio::test]
307        async fn test1_debezium_json_parser_update() {
308            //     "before": {
309            //       "id": 102,
310            //       "name": "car battery",
311            //       "description": "12V car battery",
312            //       "weight": 8.1
313            //     },
314            //     "after": {
315            //       "id": 102,
316            //       "name": "car battery",
317            //       "description": "24V car battery",
318            //       "weight": 9.1
319            //     },
320            let input = vec![
321                // data with payload field
322                br#"{"payload":{"before":{"id":102,"name":"car battery","description":"12V car battery","weight":8.1},"after":{"id":102,"name":"car battery","description":"24V car battery","weight":9.1},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551901000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1382,"row":0,"thread":null,"query":null},"op":"u","ts_ms":1639551901165,"transaction":null}}"#.to_vec(),
323                // data without payload field
324                br#"{"before":{"id":102,"name":"car battery","description":"12V car battery","weight":8.1},"after":{"id":102,"name":"car battery","description":"24V car battery","weight":9.1},"source":{"version":"1.7.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1639551901000,"snapshot":"false","db":"inventory","sequence":null,"table":"products","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1382,"row":0,"thread":null,"query":null},"op":"u","ts_ms":1639551901165,"transaction":null}"#.to_vec()];
325
326            let columns = get_test1_columns();
327
328            for data in input {
329                let parser = build_parser(columns.clone()).await;
330                let [(op, row)]: [_; 1] = parse_one(parser, columns.clone(), data)
331                    .await
332                    .try_into()
333                    .unwrap();
334
335                assert_eq!(op, Op::Insert);
336
337                assert!(row[0].eq(&Some(ScalarImpl::Int32(102))));
338                assert!(row[1].eq(&Some(ScalarImpl::Utf8("car battery".into()))));
339                assert!(row[2].eq(&Some(ScalarImpl::Utf8("24V car battery".into()))));
340                assert!(row[3].eq(&Some(ScalarImpl::Float64(9.1.into()))));
341            }
342        }
343    }
344    // test2 covers read/insert/update/delete event on the following MySQL table for debezium json:
345    // CREATE TABLE IF NOT EXISTS orders (
346    //     O_KEY BIGINT NOT NULL,
347    //     O_BOOL BOOLEAN,
348    //     O_TINY TINYINT,
349    //     O_INT INT,
350    //     O_REAL REAL,
351    //     O_DOUBLE DOUBLE,
352    //     O_DECIMAL DECIMAL(15, 2),
353    //     O_CHAR CHAR(15),
354    //     O_DATE DATE,
355    //     O_TIME TIME,
356    //     O_DATETIME DATETIME,
357    //     O_TIMESTAMP TIMESTAMP,
358    //     O_JSON JSON,
359    //     PRIMARY KEY (O_KEY));
360    // test2 also covers overflow tests on basic types
361    mod test2_mysql {
362        use super::*;
363
364        fn get_test2_columns() -> Vec<SourceColumnDesc> {
365            vec![
366                SourceColumnDesc::simple("O_KEY", DataType::Int64, ColumnId::from(0)),
367                SourceColumnDesc::simple("O_BOOL", DataType::Boolean, ColumnId::from(1)),
368                SourceColumnDesc::simple("O_TINY", DataType::Int16, ColumnId::from(2)),
369                SourceColumnDesc::simple("O_INT", DataType::Int32, ColumnId::from(3)),
370                SourceColumnDesc::simple("O_REAL", DataType::Float32, ColumnId::from(4)),
371                SourceColumnDesc::simple("O_DOUBLE", DataType::Float64, ColumnId::from(5)),
372                SourceColumnDesc::simple("O_DECIMAL", DataType::Decimal, ColumnId::from(6)),
373                SourceColumnDesc::simple("O_CHAR", DataType::Varchar, ColumnId::from(7)),
374                SourceColumnDesc::simple("O_DATE", DataType::Date, ColumnId::from(8)),
375                SourceColumnDesc::simple("O_TIME", DataType::Time, ColumnId::from(9)),
376                SourceColumnDesc::simple("O_DATETIME", DataType::Timestamp, ColumnId::from(10)),
377                SourceColumnDesc::simple("O_TIMESTAMP", DataType::Timestamptz, ColumnId::from(11)),
378                SourceColumnDesc::simple("O_JSON", DataType::Jsonb, ColumnId::from(12)),
379            ]
380        }
381
382        #[tokio::test]
383        async fn test2_debezium_json_parser_read() {
384            let data = br#"{"payload":{"before":null,"after":{"O_KEY":111,"O_BOOL":1,"O_TINY":-1,"O_INT":-1111,"O_REAL":-11.11,"O_DOUBLE":-111.11111,"O_DECIMAL":-111.11,"O_CHAR":"yes please","O_DATE":"1000-01-01","O_TIME":0,"O_DATETIME":0,"O_TIMESTAMP":"1970-01-01T00:00:01Z","O_JSON":"{\"k1\": \"v1\", \"k2\": 11}"},"source":{"version":"1.9.7.Final","connector":"mysql","name":"RW_CDC_test.orders","ts_ms":1678090651000,"snapshot":"last","db":"test","sequence":null,"table":"orders","server_id":0,"gtid":null,"file":"mysql-bin.000003","pos":951,"row":0,"thread":null,"query":null},"op":"r","ts_ms":1678090651640,"transaction":null}}"#;
385
386            let columns = get_test2_columns();
387
388            let parser = build_parser(columns.clone()).await;
389
390            let [(_op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
391                .await
392                .try_into()
393                .unwrap();
394
395            assert!(row[0].eq(&Some(ScalarImpl::Int64(111))));
396            assert!(row[1].eq(&Some(ScalarImpl::Bool(true))));
397            assert!(row[2].eq(&Some(ScalarImpl::Int16(-1))));
398            assert!(row[3].eq(&Some(ScalarImpl::Int32(-1111))));
399            assert!(row[4].eq(&Some(ScalarImpl::Float32((-11.11).into()))));
400            assert!(row[5].eq(&Some(ScalarImpl::Float64((-111.11111).into()))));
401            assert!(row[6].eq(&Some(ScalarImpl::Decimal("-111.11".parse().unwrap()))));
402            assert!(row[7].eq(&Some(ScalarImpl::Utf8("yes please".into()))));
403            assert!(row[8].eq(&Some(ScalarImpl::Date(Date::new(
404                NaiveDate::from_ymd_opt(1000, 1, 1).unwrap()
405            )))));
406            assert!(row[9].eq(&Some(ScalarImpl::Time(Time::new(
407                NaiveTime::from_hms_micro_opt(0, 0, 0, 0).unwrap()
408            )))));
409            assert!(row[10].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
410                "1970-01-01T00:00:00".parse().unwrap()
411            )))));
412            assert!(row[11].eq(&Some(ScalarImpl::Timestamptz(
413                "1970-01-01T00:00:01Z".parse().unwrap()
414            ))));
415            assert_json_eq(&row[12], "{\"k1\": \"v1\", \"k2\": 11}");
416        }
417
418        #[tokio::test]
419        async fn test2_debezium_json_parser_insert() {
420            let data = br#"{"payload":{"before":null,"after":{"O_KEY":111,"O_BOOL":1,"O_TINY":-1,"O_INT":-1111,"O_REAL":-11.11,"O_DOUBLE":-111.11111,"O_DECIMAL":-111.11,"O_CHAR":"yes please","O_DATE":"1000-01-01","O_TIME":0,"O_DATETIME":0,"O_TIMESTAMP":"1970-01-01T00:00:01Z","O_JSON":"{\"k1\": \"v1\", \"k2\": 11}"},"source":{"version":"1.9.7.Final","connector":"mysql","name":"RW_CDC_test.orders","ts_ms":1678088861000,"snapshot":"false","db":"test","sequence":null,"table":"orders","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":789,"row":0,"thread":4,"query":null},"op":"c","ts_ms":1678088861249,"transaction":null}}"#;
421
422            let columns = get_test2_columns();
423            let parser = build_parser(columns.clone()).await;
424            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
425                .await
426                .try_into()
427                .unwrap();
428            assert_eq!(op, Op::Insert);
429
430            assert!(row[0].eq(&Some(ScalarImpl::Int64(111))));
431            assert!(row[1].eq(&Some(ScalarImpl::Bool(true))));
432            assert!(row[2].eq(&Some(ScalarImpl::Int16(-1))));
433            assert!(row[3].eq(&Some(ScalarImpl::Int32(-1111))));
434            assert!(row[4].eq(&Some(ScalarImpl::Float32((-11.11).into()))));
435            assert!(row[5].eq(&Some(ScalarImpl::Float64((-111.11111).into()))));
436            assert!(row[6].eq(&Some(ScalarImpl::Decimal("-111.11".parse().unwrap()))));
437            assert!(row[7].eq(&Some(ScalarImpl::Utf8("yes please".into()))));
438            assert!(row[8].eq(&Some(ScalarImpl::Date(Date::new(
439                NaiveDate::from_ymd_opt(1000, 1, 1).unwrap()
440            )))));
441            assert!(row[9].eq(&Some(ScalarImpl::Time(Time::new(
442                NaiveTime::from_hms_micro_opt(0, 0, 0, 0).unwrap()
443            )))));
444            assert!(row[10].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
445                "1970-01-01T00:00:00".parse().unwrap()
446            )))));
447            assert!(row[11].eq(&Some(ScalarImpl::Timestamptz(
448                "1970-01-01T00:00:01Z".parse().unwrap()
449            ))));
450            assert_json_eq(&row[12], "{\"k1\": \"v1\", \"k2\": 11}");
451        }
452
453        #[tokio::test]
454        async fn test2_debezium_json_parser_delete() {
455            let data = br#"{"payload":{"before":{"O_KEY":111,"O_BOOL":0,"O_TINY":3,"O_INT":3333,"O_REAL":33.33,"O_DOUBLE":333.33333,"O_DECIMAL":333.33,"O_CHAR":"no thanks","O_DATE":"9999-12-31","O_TIME":86399000000,"O_DATETIME":99999999999000,"O_TIMESTAMP":"2038-01-09T03:14:07Z","O_JSON":"{\"k1\":\"v1_updated\",\"k2\":33}"},"after":null,"source":{"version":"1.9.7.Final","connector":"mysql","name":"RW_CDC_test.orders","ts_ms":1678090653000,"snapshot":"false","db":"test","sequence":null,"table":"orders","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1643,"row":0,"thread":4,"query":null},"op":"d","ts_ms":1678090653611,"transaction":null}}"#;
456
457            let columns = get_test2_columns();
458            let parser = build_parser(columns.clone()).await;
459            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
460                .await
461                .try_into()
462                .unwrap();
463
464            assert_eq!(op, Op::Delete);
465
466            assert!(row[0].eq(&Some(ScalarImpl::Int64(111))));
467            assert!(row[1].eq(&Some(ScalarImpl::Bool(false))));
468            assert!(row[2].eq(&Some(ScalarImpl::Int16(3))));
469            assert!(row[3].eq(&Some(ScalarImpl::Int32(3333))));
470            assert!(row[4].eq(&Some(ScalarImpl::Float32((33.33).into()))));
471            assert!(row[5].eq(&Some(ScalarImpl::Float64((333.33333).into()))));
472            assert!(row[6].eq(&Some(ScalarImpl::Decimal("333.33".parse().unwrap()))));
473            assert!(row[7].eq(&Some(ScalarImpl::Utf8("no thanks".into()))));
474            assert!(row[8].eq(&Some(ScalarImpl::Date(Date::new(
475                NaiveDate::from_ymd_opt(9999, 12, 31).unwrap()
476            )))));
477            assert!(row[9].eq(&Some(ScalarImpl::Time(Time::new(
478                NaiveTime::from_hms_micro_opt(23, 59, 59, 0).unwrap()
479            )))));
480            assert!(row[10].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
481                "5138-11-16T09:46:39".parse().unwrap()
482            )))));
483            assert!(row[11].eq(&Some(ScalarImpl::Timestamptz(
484                "2038-01-09T03:14:07Z".parse().unwrap()
485            ))));
486            assert_json_eq(&row[12], "{\"k1\":\"v1_updated\",\"k2\":33}");
487        }
488
489        #[tokio::test]
490        async fn test2_debezium_json_parser_update() {
491            let data = br#"{"payload":{"before":{"O_KEY":111,"O_BOOL":1,"O_TINY":-1,"O_INT":-1111,"O_REAL":-11.11,"O_DOUBLE":-111.11111,"O_DECIMAL":-111.11,"O_CHAR":"yes please","O_DATE":"1000-01-01","O_TIME":0,"O_DATETIME":0,"O_TIMESTAMP":"1970-01-01T00:00:01Z","O_JSON":"{\"k1\": \"v1\", \"k2\": 11}"},"after":{"O_KEY":111,"O_BOOL":0,"O_TINY":3,"O_INT":3333,"O_REAL":33.33,"O_DOUBLE":333.33333,"O_DECIMAL":333.33,"O_CHAR":"no thanks","O_DATE":"9999-12-31","O_TIME":86399000000,"O_DATETIME":99999999999000,"O_TIMESTAMP":"2038-01-09T03:14:07Z","O_JSON":"{\"k1\": \"v1_updated\", \"k2\": 33}"},"source":{"version":"1.9.7.Final","connector":"mysql","name":"RW_CDC_test.orders","ts_ms":1678089331000,"snapshot":"false","db":"test","sequence":null,"table":"orders","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":1168,"row":0,"thread":4,"query":null},"op":"u","ts_ms":1678089331464,"transaction":null}}"#;
492
493            let columns = get_test2_columns();
494
495            let parser = build_parser(columns.clone()).await;
496            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
497                .await
498                .try_into()
499                .unwrap();
500
501            assert_eq!(op, Op::Insert);
502
503            assert!(row[0].eq(&Some(ScalarImpl::Int64(111))));
504            assert!(row[1].eq(&Some(ScalarImpl::Bool(false))));
505            assert!(row[2].eq(&Some(ScalarImpl::Int16(3))));
506            assert!(row[3].eq(&Some(ScalarImpl::Int32(3333))));
507            assert!(row[4].eq(&Some(ScalarImpl::Float32((33.33).into()))));
508            assert!(row[5].eq(&Some(ScalarImpl::Float64((333.33333).into()))));
509            assert!(row[6].eq(&Some(ScalarImpl::Decimal("333.33".parse().unwrap()))));
510            assert!(row[7].eq(&Some(ScalarImpl::Utf8("no thanks".into()))));
511            assert!(row[8].eq(&Some(ScalarImpl::Date(Date::new(
512                NaiveDate::from_ymd_opt(9999, 12, 31).unwrap()
513            )))));
514            assert!(row[9].eq(&Some(ScalarImpl::Time(Time::new(
515                NaiveTime::from_hms_micro_opt(23, 59, 59, 0).unwrap()
516            )))));
517            assert!(row[10].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
518                "5138-11-16T09:46:39".parse().unwrap()
519            )))));
520            assert!(row[11].eq(&Some(ScalarImpl::Timestamptz(
521                "2038-01-09T03:14:07Z".parse().unwrap()
522            ))));
523            assert_json_eq(&row[12], "{\"k1\": \"v1_updated\", \"k2\": 33}");
524        }
525
526        #[cfg(not(madsim))] // Traced test does not work with madsim
527        #[tokio::test]
528        #[tracing_test::traced_test]
529        async fn test2_debezium_json_parser_overflow() {
530            let columns = vec![
531                SourceColumnDesc::simple("O_KEY", DataType::Int64, ColumnId::from(0)),
532                SourceColumnDesc::simple("O_BOOL", DataType::Boolean, ColumnId::from(1)),
533                SourceColumnDesc::simple("O_TINY", DataType::Int16, ColumnId::from(2)),
534                SourceColumnDesc::simple("O_INT", DataType::Int32, ColumnId::from(3)),
535                SourceColumnDesc::simple("O_REAL", DataType::Float32, ColumnId::from(4)),
536                SourceColumnDesc::simple("O_DOUBLE", DataType::Float64, ColumnId::from(5)),
537            ];
538            let mut parser = build_parser(columns.clone()).await;
539
540            let mut dummy_builder =
541                SourceStreamChunkBuilder::new(columns, SourceCtrlOpts::for_test());
542
543            let normal_values = ["111", "1", "33", "444", "555.0", "666.0"];
544            let overflow_values = [
545                "9223372036854775808",
546                "2",
547                "32768",
548                "2147483648",
549                "3.80282347E38",
550                "1.797695E308",
551            ];
552
553            for i in 0..6 {
554                let mut values = normal_values;
555                values[i] = overflow_values[i];
556                let data = format!(
557                    r#"{{"payload":{{"before":null,"after":{{"O_KEY":{},"O_BOOL":{},"O_TINY":{},"O_INT":{},"O_REAL":{},"O_DOUBLE":{}}},"source":{{"version":"1.9.7.Final","connector":"mysql","name":"RW_CDC_test.orders","ts_ms":1678158055000,"snapshot":"false","db":"test","sequence":null,"table":"orders","server_id":223344,"gtid":null,"file":"mysql-bin.000003","pos":637,"row":0,"thread":4,"query":null}},"op":"c","ts_ms":1678158055464,"transaction":null}}}}"#,
558                    values[0], values[1], values[2], values[3], values[4], values[5]
559                ).as_bytes().to_vec();
560
561                let res = parser
562                    .parse_inner(None, Some(data), dummy_builder.row_writer())
563                    .await;
564                if i < 5 {
565                    // For other overflow, the parsing succeeds but the type conversion fails
566                    // The errors are ignored and logged.
567                    res.unwrap();
568                    assert!(logs_contain("expected type"), "{i}");
569                } else {
570                    // For f64 overflow, the parsing fails
571                    let e = res.unwrap_err();
572                    assert!(e.to_report_string().contains("InvalidNumber"), "{i}: {e}");
573                }
574            }
575        }
576    }
577
578    // postgres-specific data-type mapping tests
579    mod test3_postgres {
580        use risingwave_pb::plan_common::AdditionalColumn;
581
582        use super::*;
583        use crate::connector_common::postgres::postgres_point_type;
584        use crate::source::SourceColumnType;
585
586        // schema for temporal-type test
587        fn get_temporal_test_columns() -> Vec<SourceColumnDesc> {
588            vec![
589                SourceColumnDesc::simple("o_key", DataType::Int32, ColumnId::from(0)),
590                SourceColumnDesc::simple("o_time_0", DataType::Time, ColumnId::from(1)),
591                SourceColumnDesc::simple("o_time_6", DataType::Time, ColumnId::from(2)),
592                SourceColumnDesc::simple("o_timez_0", DataType::Time, ColumnId::from(3)),
593                SourceColumnDesc::simple("o_timez_6", DataType::Time, ColumnId::from(4)),
594                SourceColumnDesc::simple("o_timestamp_0", DataType::Timestamp, ColumnId::from(5)),
595                SourceColumnDesc::simple("o_timestamp_6", DataType::Timestamp, ColumnId::from(6)),
596                SourceColumnDesc::simple(
597                    "o_timestampz_0",
598                    DataType::Timestamptz,
599                    ColumnId::from(7),
600                ),
601                SourceColumnDesc::simple(
602                    "o_timestampz_6",
603                    DataType::Timestamptz,
604                    ColumnId::from(8),
605                ),
606                SourceColumnDesc::simple("o_interval", DataType::Interval, ColumnId::from(9)),
607                SourceColumnDesc::simple("o_date", DataType::Date, ColumnId::from(10)),
608            ]
609        }
610
611        // schema for numeric-type test
612        fn get_numeric_test_columns() -> Vec<SourceColumnDesc> {
613            vec![
614                SourceColumnDesc::simple("o_key", DataType::Int32, ColumnId::from(0)),
615                SourceColumnDesc::simple("o_smallint", DataType::Int16, ColumnId::from(1)),
616                SourceColumnDesc::simple("o_integer", DataType::Int32, ColumnId::from(2)),
617                SourceColumnDesc::simple("o_bigint", DataType::Int64, ColumnId::from(3)),
618                SourceColumnDesc::simple("o_real", DataType::Float32, ColumnId::from(4)),
619                SourceColumnDesc::simple("o_double", DataType::Float64, ColumnId::from(5)),
620                SourceColumnDesc::simple("o_numeric", DataType::Decimal, ColumnId::from(6)),
621                SourceColumnDesc::simple("o_numeric_6_3", DataType::Decimal, ColumnId::from(7)),
622                SourceColumnDesc::simple("o_money", DataType::Decimal, ColumnId::from(8)),
623            ]
624        }
625
626        // schema for the remaining types
627        fn get_other_types_test_columns() -> Vec<SourceColumnDesc> {
628            vec![
629                SourceColumnDesc::simple("o_key", DataType::Int32, ColumnId::from(0)),
630                SourceColumnDesc::simple("o_boolean", DataType::Boolean, ColumnId::from(1)),
631                SourceColumnDesc::simple("o_bit", DataType::Boolean, ColumnId::from(2)),
632                SourceColumnDesc::simple("o_bytea", DataType::Bytea, ColumnId::from(3)),
633                SourceColumnDesc::simple("o_json", DataType::Jsonb, ColumnId::from(4)),
634                SourceColumnDesc::simple("o_xml", DataType::Varchar, ColumnId::from(5)),
635                SourceColumnDesc::simple("o_uuid", DataType::Varchar, ColumnId::from(6)),
636                SourceColumnDesc {
637                    name: "o_point".to_owned(),
638                    data_type: postgres_point_type(),
639                    column_id: 7.into(),
640                    column_type: SourceColumnType::Normal,
641                    is_pk: false,
642                    is_hidden_addition_col: false,
643                    additional_column: AdditionalColumn { column_type: None },
644                },
645                SourceColumnDesc::simple("o_enum", DataType::Varchar, ColumnId::from(8)),
646                SourceColumnDesc::simple("o_char", DataType::Varchar, ColumnId::from(9)),
647                SourceColumnDesc::simple("o_varchar", DataType::Varchar, ColumnId::from(10)),
648                SourceColumnDesc::simple("o_character", DataType::Varchar, ColumnId::from(11)),
649                SourceColumnDesc::simple(
650                    "o_character_varying",
651                    DataType::Varchar,
652                    ColumnId::from(12),
653                ),
654            ]
655        }
656
657        #[tokio::test]
658        async fn test_temporal_types() {
659            // this test includes all supported temporal types, with the schema
660            // CREATE TABLE orders (
661            //     o_key integer,
662            //     o_time_0 time(0),
663            //     o_time_6 time(6),
664            //     o_timez_0 time(0) with time zone,
665            //     o_timez_6 time(6) with time zone,
666            //     o_timestamp_0 timestamp(0),
667            //     o_timestamp_6 timestamp(6),
668            //     o_timestampz_0 timestamp(0) with time zone,
669            //     o_timestampz_6 timestamp(6) with time zone,
670            //     o_interval interval,
671            //     o_date date,
672            //     PRIMARY KEY (o_key)
673            // );
674            // this test covers an insert event on the table above
675            let data = br#"{"payload":{"before":null,"after":{"o_key":0,"o_time_0":40271000000,"o_time_6":40271000010,"o_timez_0":"11:11:11Z","o_timez_6":"11:11:11.00001Z","o_timestamp_0":1321009871000,"o_timestamp_6":1321009871123456,"o_timestampz_0":"2011-11-11T03:11:11Z","o_timestampz_6":"2011-11-11T03:11:11.123456Z","o_interval":"P1Y2M3DT4H5M6.78S","o_date":"1999-09-09"},"source":{"version":"1.9.7.Final","connector":"postgresql","name":"RW_CDC_localhost.test.orders","ts_ms":1684733351963,"snapshot":"last","db":"test","sequence":"[null,\"26505352\"]","schema":"public","table":"orders","txId":729,"lsn":26505352,"xmin":null},"op":"r","ts_ms":1684733352110,"transaction":null}}"#;
676            let columns = get_temporal_test_columns();
677            let parser = build_parser(columns.clone()).await;
678            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
679                .await
680                .try_into()
681                .unwrap();
682            assert_eq!(op, Op::Insert);
683            assert!(row[0].eq(&Some(ScalarImpl::Int32(0))));
684            assert!(row[1].eq(&Some(ScalarImpl::Time(Time::new(
685                NaiveTime::from_hms_micro_opt(11, 11, 11, 0).unwrap()
686            )))));
687            assert!(row[2].eq(&Some(ScalarImpl::Time(Time::new(
688                NaiveTime::from_hms_micro_opt(11, 11, 11, 10).unwrap()
689            )))));
690            assert!(row[3].eq(&Some(ScalarImpl::Time(Time::new(
691                NaiveTime::from_hms_micro_opt(11, 11, 11, 0).unwrap()
692            )))));
693            assert!(row[4].eq(&Some(ScalarImpl::Time(Time::new(
694                NaiveTime::from_hms_micro_opt(11, 11, 11, 10).unwrap()
695            )))));
696            assert!(row[5].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
697                "2011-11-11T11:11:11".parse().unwrap()
698            )))));
699            assert!(row[6].eq(&Some(ScalarImpl::Timestamp(Timestamp::new(
700                "2011-11-11T11:11:11.123456".parse().unwrap()
701            )))));
702            assert!(
703                row[9].eq(&Some(ScalarImpl::Interval(Interval::from_month_day_usec(
704                    14,
705                    3,
706                    14706780000
707                ))))
708            );
709            assert!(row[10].eq(&Some(ScalarImpl::Date(Date::new(
710                NaiveDate::from_ymd_opt(1999, 9, 9).unwrap()
711            )))));
712        }
713
714        #[tokio::test]
715        async fn test_numeric_types() {
716            // this test includes all supported numeric types, with the schema
717            // CREATE TABLE orders (
718            //     o_key integer,
719            //     o_smallint smallint,
720            //     o_integer integer,
721            //     o_bigint bigint,
722            //     o_real real,
723            //     o_double double precision,
724            //     o_numeric numeric,
725            //     o_numeric_6_3 numeric(6,3),
726            //     o_money money,
727            //     PRIMARY KEY (o_key)
728            // );
729            // this test covers an insert event on the table above
730            let data = br#"{"payload":{"before":null,"after":{"o_key":0,"o_smallint":32767,"o_integer":2147483647,"o_bigint":9223372036854775807,"o_real":9.999,"o_double":9.999999,"o_numeric":123456.789,"o_numeric_6_3":123.456,"o_money":123.12},"source":{"version":"1.9.7.Final","connector":"postgresql","name":"RW_CDC_localhost.test.orders","ts_ms":1684404343201,"snapshot":"last","db":"test","sequence":"[null,\"26519216\"]","schema":"public","table":"orders","txId":729,"lsn":26519216,"xmin":null},"op":"r","ts_ms":1684404343349,"transaction":null}}"#;
731            let columns = get_numeric_test_columns();
732            let parser = build_parser(columns.clone()).await;
733            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
734                .await
735                .try_into()
736                .unwrap();
737            assert_eq!(op, Op::Insert);
738            assert!(row[0].eq(&Some(ScalarImpl::Int32(0))));
739            assert!(row[1].eq(&Some(ScalarImpl::Int16(32767))));
740            assert!(row[2].eq(&Some(ScalarImpl::Int32(2147483647))));
741            assert!(row[3].eq(&Some(ScalarImpl::Int64(9223372036854775807))));
742            assert!(row[4].eq(&Some(ScalarImpl::Float32((9.999).into()))));
743            assert!(row[5].eq(&Some(ScalarImpl::Float64((9.999999).into()))));
744            assert!(row[6].eq(&Some(ScalarImpl::Decimal("123456.7890".parse().unwrap()))));
745            assert!(row[7].eq(&Some(ScalarImpl::Decimal("123.456".parse().unwrap()))));
746            assert!(row[8].eq(&Some(ScalarImpl::Decimal("123.12".parse().unwrap()))));
747        }
748
749        #[tokio::test]
750        async fn test_other_types() {
751            // this test includes the remaining types, with the schema
752            // CREATE TABLE orders (
753            //     o_key integer,
754            //     o_boolean boolean,
755            //     o_bit bit,
756            //     o_bytea bytea,
757            //     o_json jsonb,
758            //     o_xml xml,
759            //     o_uuid uuid,
760            //     o_point point,
761            //     o_enum bear,
762            //     o_char char,
763            //     o_varchar varchar,
764            //     o_character character,
765            //     o_character_varying character varying,
766            //     PRIMARY KEY (o_key)
767            //  );
768            // this test covers an insert event on the table above
769            let data = br#"{"payload":{"before":null,"after":{"o_key":1,"o_boolean":false,"o_bit":true,"o_bytea":"ASNFZ4mrze8=","o_json":"{\"k1\": \"v1\", \"k2\": 11}","o_xml":"<!--hahaha-->","o_uuid":"60f14fe2-f857-404a-b586-3b5375b3259f","o_point":{"x":1.0,"y":2.0,"wkb":"AQEAAAAAAAAAAADwPwAAAAAAAABA","srid":null},"o_enum":"polar","o_char":"h","o_varchar":"ha","o_character":"h","o_character_varying":"hahaha"},"source":{"version":"1.9.7.Final","connector":"postgresql","name":"RW_CDC_localhost.test.orders","ts_ms":1684743927178,"snapshot":"last","db":"test","sequence":"[null,\"26524528\"]","schema":"public","table":"orders","txId":730,"lsn":26524528,"xmin":null},"op":"r","ts_ms":1684743927343,"transaction":null}}"#;
770            let columns = get_other_types_test_columns();
771            let parser = build_parser(columns.clone()).await;
772            let [(op, row)]: [_; 1] = parse_one(parser, columns, data.to_vec())
773                .await
774                .try_into()
775                .unwrap();
776            assert_eq!(op, Op::Insert);
777            assert!(row[0].eq(&Some(ScalarImpl::Int32(1))));
778            assert!(row[1].eq(&Some(ScalarImpl::Bool(false))));
779            assert!(row[2].eq(&Some(ScalarImpl::Bool(true))));
780            assert!(row[3].eq(&Some(ScalarImpl::Bytea(Box::new([
781                u8::from_str_radix("01", 16).unwrap(),
782                u8::from_str_radix("23", 16).unwrap(),
783                u8::from_str_radix("45", 16).unwrap(),
784                u8::from_str_radix("67", 16).unwrap(),
785                u8::from_str_radix("89", 16).unwrap(),
786                u8::from_str_radix("AB", 16).unwrap(),
787                u8::from_str_radix("CD", 16).unwrap(),
788                u8::from_str_radix("EF", 16).unwrap()
789            ])))));
790            assert_json_eq(&row[4], "{\"k1\": \"v1\", \"k2\": 11}");
791            assert!(row[5].eq(&Some(ScalarImpl::Utf8("<!--hahaha-->".into()))));
792            assert!(row[6].eq(&Some(ScalarImpl::Utf8(
793                "60f14fe2-f857-404a-b586-3b5375b3259f".into()
794            ))));
795            assert!(row[7].eq(&Some(ScalarImpl::Struct(StructValue::new(vec![
796                Some(ScalarImpl::Float64(1.into())),
797                Some(ScalarImpl::Float64(2.into()))
798            ])))));
799            assert!(row[8].eq(&Some(ScalarImpl::Utf8("polar".into()))));
800            assert!(row[9].eq(&Some(ScalarImpl::Utf8("h".into()))));
801            assert!(row[10].eq(&Some(ScalarImpl::Utf8("ha".into()))));
802            assert!(row[11].eq(&Some(ScalarImpl::Utf8("h".into()))));
803            assert!(row[12].eq(&Some(ScalarImpl::Utf8("hahaha".into()))));
804        }
805    }
806}