Skip to content

Commit 320bfbb

Browse files
committed
add jsonValue is null judgement for issue :datavane/tis#456
1 parent 0fc1ff9 commit 320bfbb

10 files changed

Lines changed: 149 additions & 60 deletions

File tree

tis-datax/tis-datax-mongodb-plugin/src/main/java/com/qlangtech/tis/plugin/datax/mongo/MongoCMeta.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -146,7 +146,7 @@ private static BsonValue getEmbeddedValue(List<String> keys, BsonDocument doc) {
146146
while (keyIterator.hasNext()) {
147147
key = keyIterator.next();
148148
value = doc.get(key);
149-
if (value == null) {
149+
if (value == null || value.isNull()) {
150150
return null;
151151
}
152152
if (value.isDocument()) {

tis-incr/pom.xml

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@
4040
<!--https://github.com/ververica/flink-cdc-connectors-->
4141
<!-- <flink.cdc.version>2.1.0</flink.cdc.version>-->
4242

43-
<chunjun.version>1.12.5-${revision}</chunjun.version>
43+
<chunjun.version>1.12.5-${project.version}</chunjun.version>
4444
<debezium-connector.version>1.9.8.Final</debezium-connector.version>
4545
</properties>
4646

@@ -60,6 +60,7 @@
6060
<module>tis-flink-dependency</module>
6161
<module>tis-flink-extends</module>
6262
<module>tis-flink-extends-plugin</module>
63+
6364
<module>tis-flink-cdc-common</module>
6465
<module>tis-flink-cdc-mysql-plugin</module>
6566
<module>tis-flink-cdc-postgresql-plugin</module>
@@ -100,7 +101,6 @@
100101
<groupId>com.qlangtech.tis.plugins</groupId>
101102
<artifactId>tis-flink-cdc-common</artifactId>
102103
<version>${project.version}</version>
103-
<scope>provided</scope>
104104
<exclusions>
105105
<exclusion>
106106
<groupId>io.debezium</groupId>
@@ -200,6 +200,7 @@
200200
<module>tis-flink-cdc-postgresql-shade-4-debezium-connector-postgresql</module>
201201
<module>tis-flink-cdc-mysql-shade-4-debezium-connector-mysql</module>
202202
<module>tis-flink-cdc-kingbase-shade-4-debezium-connector-postgresql</module>
203+
<module>tis-flink-cdc-common-plugin-shade-4-debezium</module>
203204
<!-- <module>tis-flink-cdc-kingbase-shade-plugin-extends</module>-->
204205
</modules>
205206
</profile>
Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,104 @@
1+
<?xml version="1.0" encoding="UTF-8"?>
2+
<!--~
3+
~ Licensed to the Apache Software Foundation (ASF) under one
4+
~ or more contributor license agreements. See the NOTICE file
5+
~ distributed with this work for additional information
6+
~ regarding copyright ownership. The ASF licenses this file
7+
~ to you under the Apache License, Version 2.0 (the
8+
~ "License"); you may not use this file except in compliance
9+
~ with the License. You may obtain a copy of the License at
10+
~
11+
~ http://www.apache.org/licenses/LICENSE-2.0
12+
~
13+
~ Unless required by applicable law or agreed to in writing, software
14+
~ distributed under the License is distributed on an "AS IS" BASIS,
15+
~ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16+
~ See the License for the specific language governing permissions and
17+
~ limitations under the License.
18+
-->
19+
20+
<project xmlns="http://maven.apache.org/POM/4.0.0"
21+
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
22+
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
23+
<modelVersion>4.0.0</modelVersion>
24+
<parent>
25+
<groupId>com.qlangtech.tis.plugins</groupId>
26+
<artifactId>tis-incr</artifactId>
27+
<version>${revision}</version>
28+
<relativePath>../pom.xml</relativePath>
29+
</parent>
30+
31+
<artifactId>tis-flink-cdc-common-plugin-shade-4-debuezium</artifactId>
32+
33+
<properties>
34+
35+
</properties>
36+
37+
<dependencies>
38+
<dependency>
39+
<groupId>org.apache.flink</groupId>
40+
<artifactId>flink-connector-debezium</artifactId>
41+
<version>${flink.cdc.version}</version>
42+
</dependency>
43+
</dependencies>
44+
45+
<build>
46+
<plugins>
47+
48+
<!-- Build uber jar -->
49+
<plugin>
50+
<groupId>org.apache.maven.plugins</groupId>
51+
<artifactId>maven-shade-plugin</artifactId>
52+
<executions>
53+
<execution>
54+
<phase>package</phase>
55+
<goals>
56+
<goal>shade</goal>
57+
</goals>
58+
<configuration>
59+
<createDependencyReducedPom>false</createDependencyReducedPom>
60+
<shadedArtifactAttached>false</shadedArtifactAttached>
61+
<!-- <finalName>${project.artifactId}-dist-${project.version}</finalName>-->
62+
<artifactSet>
63+
<excludes>
64+
<exclude>log4j:log4j</exclude>
65+
<exclude>org.slf4j:slf4j-api</exclude>
66+
<exclude>org.codehaus.groovy:groovy-all</exclude>
67+
</excludes>
68+
</artifactSet>
69+
<filters>
70+
<filter>
71+
<artifact>*:*</artifact>
72+
<excludes>
73+
<exclude>META-INF/*.SF</exclude>
74+
<exclude>META-INF/*.DSA</exclude>
75+
<exclude>META-INF/*.RSA</exclude>
76+
</excludes>
77+
</filter>
78+
<filter>
79+
<artifact>io.debezium:debezium-embedded</artifact>
80+
<excludes>
81+
<!--https://github.com/datavane/tis/issues/454-->
82+
<!--have been rewrite in flink-connector-debezium-->
83+
<exclude>io/debezium/embedded/EmbeddedEngineChangeEvent.class</exclude>
84+
</excludes>
85+
</filter>
86+
<filter>
87+
<artifact>io.debezium:debezium-core</artifact>
88+
<excludes>
89+
<exclude>io/debezium/relational/HistorizedRelationalDatabaseConnectorConfig*</exclude>
90+
<exclude>io/debezium/relational/RelationalChangeRecordEmitter*</exclude>
91+
<exclude>io/debezium/relational/RelationalTableFilters*</exclude>
92+
</excludes>
93+
</filter>
94+
</filters>
95+
</configuration>
96+
</execution>
97+
98+
</executions>
99+
</plugin>
100+
101+
</plugins>
102+
</build>
103+
104+
</project>

tis-incr/tis-flink-cdc-common/pom.xml

Lines changed: 21 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -30,33 +30,39 @@
3030

3131
<groupId>com.qlangtech.tis.plugins</groupId>
3232
<artifactId>tis-flink-cdc-common</artifactId>
33-
<!-- <packaging>tpi</packaging>-->
33+
<packaging>tpi</packaging>
3434

3535
<dependencies>
36+
3637
<dependency>
37-
<groupId>org.apache.flink</groupId>
38-
<artifactId>flink-connector-debezium</artifactId>
39-
<version>${flink.cdc.version}</version>
40-
<!-- <scope>provided</scope>-->
41-
<!-- <exclusions>-->
42-
<!-- <exclusion>-->
43-
<!-- <groupId>io.debezium</groupId>-->
44-
<!-- <artifactId>debezium-core</artifactId>-->
45-
<!-- </exclusion>-->
46-
<!-- </exclusions>-->
38+
<groupId>com.qlangtech.tis.plugins</groupId>
39+
<artifactId>tis-flink-cdc-common-plugin-shade-4-debuezium</artifactId>
40+
<version>${project.version}</version>
41+
<exclusions>
42+
<exclusion>
43+
<groupId>*</groupId>
44+
<artifactId>*</artifactId>
45+
</exclusion>
46+
</exclusions>
4747
</dependency>
4848

49+
<!-- <dependency>-->
50+
<!-- <groupId>org.apache.flink</groupId>-->
51+
<!-- <artifactId>flink-connector-debezium</artifactId>-->
52+
<!-- <version>${flink.cdc.version}</version>-->
53+
<!-- </dependency>-->
54+
4955
<dependency>
5056
<groupId>org.apache.flink</groupId>
5157
<artifactId>flink-core</artifactId>
5258
<version>${flink.version}</version>
5359
<scope>provided</scope>
5460
</dependency>
5561

56-
<!-- <dependency>-->
57-
<!-- <groupId>com.qlangtech.tis.plugins</groupId>-->
58-
<!-- <artifactId>tis-flink-dependency</artifactId>-->
59-
<!-- </dependency>-->
62+
<!-- <dependency>-->
63+
<!-- <groupId>com.qlangtech.tis.plugins</groupId>-->
64+
<!-- <artifactId>tis-flink-dependency</artifactId>-->
65+
<!-- </dependency>-->
6066

6167
</dependencies>
6268

tis-incr/tis-flink-cdc-mysql-plugin/pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -150,9 +150,9 @@
150150
<plugin>
151151
<groupId>com.qlangtech.tis</groupId>
152152
<artifactId>maven-tpi-plugin</artifactId>
153-
<configuration>
154-
<maskClasses>org.apache.kafka.</maskClasses>
155-
</configuration>
153+
<!-- <configuration>-->
154+
<!-- <maskClasses>org.apache.kafka.</maskClasses>-->
155+
<!-- </configuration>-->
156156
</plugin>
157157
</plugins>
158158
</build>

tis-incr/tis-flink-cdc-mysql-plugin/src/test/java/com/qlangtech/tis/plugins/incr/flink/cdc/mysql/TestFlinkCDCMySQLSourceFactory.java

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -206,7 +206,7 @@ public void testBinlogConsumeWithRowTransformer() throws Exception {
206206
streamFactory.parallelism = 1;
207207
FlinkCDCMySQLSourceFactory mysqlCDCFactory = createCDCFactory();
208208
mysqlCDCFactory.startupOptions = new LatestStartupOptions();
209-
mysqlCDCFactory.timeZone = FlinkCDCMySQLSourceFactory.dftZoneId();
209+
210210

211211
CDCTestSuitParams suitParams = tabParamMap.get(tabBase);
212212
CUDCDCTestSuit cdcTestSuit = new CUDCDCTestSuit(suitParams) {
@@ -340,7 +340,7 @@ public void testFullTypesConsume() throws Exception {
340340
streamFactory.parallelism = 1;
341341
FlinkCDCMySQLSourceFactory mysqlCDCFactory = createCDCFactory();
342342
mysqlCDCFactory.startupOptions = new LatestStartupOptions();
343-
mysqlCDCFactory.timeZone = FlinkCDCMySQLSourceFactory.dftZoneId();
343+
// mysqlCDCFactory.timeZone = FlinkCDCMySQLSourceFactory.dftZoneId();
344344
// final String tabName = "base";
345345
CDCTestSuitParams suitParams = tabParamMap.get(fullTypes);
346346
Assert.assertNotNull(suitParams);
@@ -463,7 +463,7 @@ protected IResultRows createConsumerHandle(BasicDataXRdbmsReader dataxReader, St
463463

464464
protected FlinkCDCMySQLSourceFactory createCDCFactory() {
465465
FlinkCDCMySQLSourceFactory mySQLSourceFactory = new FlinkCDCMySQLSourceFactory();
466-
mySQLSourceFactory.timeZone = FlinkCDCMySQLSourceFactory.dftZoneId();
466+
// mySQLSourceFactory.timeZone = FlinkCDCMySQLSourceFactory.dftZoneId();
467467
return mySQLSourceFactory;
468468
}
469469

@@ -703,6 +703,11 @@ public Void smallIntType(DataType dataType) {
703703
tinyIntType(dataType);
704704
return null;
705705
}
706+
707+
@Override
708+
public Void boolType(DataType dataType) {
709+
throw new UnsupportedOperationException();
710+
}
706711
});
707712
});
708713

tis-incr/tis-flink-cdc-mysql-plugin/src/test/java/com/qlangtech/tis/plugins/incr/flink/cdc/mysql/TestTime.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ public class TestTime {
3636
@Test
3737
public void testURLs() throws Exception {
3838
Enumeration<URL> resources = TestTime.class.getClassLoader().getResources(
39-
"io/debezium/relational/RelationalTableFilters.class");
39+
"org/apache/kafka/connect/errors/ConnectException.class");
4040
while (resources.hasMoreElements()) {
4141
System.out.println(resources.nextElement());
4242
}

tis-incr/tis-flink-cdc-mysql-shade-4-debezium-connector-mysql/dependency-reduced-pom.xml

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -36,14 +36,6 @@
3636
<exclude>io/debezium/connector/mysql/MySqlConnection*</exclude>
3737
</excludes>
3838
</filter>
39-
<filter>
40-
<artifact>io.debezium:debezium-core</artifact>
41-
<excludes>
42-
<exclude>io/debezium/relational/HistorizedRelationalDatabaseConnectorConfig*</exclude>
43-
<exclude>io/debezium/relational/RelationalChangeRecordEmitter*</exclude>
44-
<exclude>io/debezium/relational/RelationalTableFilters*</exclude>
45-
</excludes>
46-
</filter>
4739
</filters>
4840
<relocations />
4941
<artifactSet>

tis-incr/tis-flink-extends-plugin/pom.xml

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,13 @@
3737
<groupId>com.qlangtech.tis.plugins</groupId>
3838
<artifactId>tis-flink-extends</artifactId>
3939
<version>${project.version}</version>
40+
<!-- <classifier>shaded</classifier>-->
41+
<!-- <exclusions>-->
42+
<!-- <exclusion>-->
43+
<!-- <groupId>*</groupId>-->
44+
<!-- <artifactId>*</artifactId>-->
45+
<!-- </exclusion>-->
46+
<!-- </exclusions>-->
4047
<exclusions>
4148
<exclusion>
4249
<groupId>com.qlangtech.tis</groupId>

tis-incr/tis-flink-extends/pom.xml

Lines changed: 1 addition & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -35,17 +35,7 @@
3535

3636
<dependencies>
3737

38-
<dependency>
39-
<groupId>com.qlangtech.tis.plugins</groupId>
40-
<artifactId>tis-flink-cdc-common</artifactId>
41-
<scope>compile</scope>
42-
</dependency>
43-
<dependency>
44-
<groupId>io.debezium</groupId>
45-
<artifactId>debezium-core</artifactId>
46-
<!--当flink-cdc 依赖的 debezium 版本变化,此版本也需要相应变化成对应的版本-->
47-
<version>${debezium-connector.version}</version>
48-
</dependency>
38+
4939

5040
<dependency>
5141
<groupId>org.apache.flink</groupId>
@@ -188,22 +178,6 @@
188178
<exclude>META-INF/*.RSA</exclude>
189179
</excludes>
190180
</filter>
191-
<filter>
192-
<artifact>io.debezium:debezium-embedded</artifact>
193-
<excludes>
194-
<!--https://github.com/datavane/tis/issues/454-->
195-
<!--have been rewrite in flink-connector-debezium-->
196-
<exclude>io/debezium/embedded/EmbeddedEngineChangeEvent.class</exclude>
197-
</excludes>
198-
</filter>
199-
<filter>
200-
<artifact>io.debezium:debezium-core</artifact>
201-
<excludes>
202-
<exclude>io/debezium/relational/HistorizedRelationalDatabaseConnectorConfig*</exclude>
203-
<exclude>io/debezium/relational/RelationalChangeRecordEmitter*</exclude>
204-
<exclude>io/debezium/relational/RelationalTableFilters*</exclude>
205-
</excludes>
206-
</filter>
207181
</filters>
208182
<!-- <transformers>-->
209183
<!-- <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">-->

0 commit comments

Comments
 (0)