Skip to content

Commit 11b25fa

Browse files
committed
Use a kafka-safe table in FlinkMainIT
Also added the custom classes to the uber JAR for proper classloading
1 parent 7af5ed9 commit 11b25fa

File tree

16 files changed

+320
-117
lines changed

16 files changed

+320
-117
lines changed

connectors/kafka-safe/pom.xml

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,72 @@
1+
<?xml version="1.0" encoding="UTF-8"?>
2+
<!--
3+
4+
Copyright © 2024 DataSQRL (contact@datasqrl.com)
5+
6+
Licensed under the Apache License, Version 2.0 (the "License");
7+
you may not use this file except in compliance with the License.
8+
You may obtain a copy of the License at
9+
10+
http://www.apache.org/licenses/LICENSE-2.0
11+
12+
Unless required by applicable law or agreed to in writing, software
13+
distributed under the License is distributed on an "AS IS" BASIS,
14+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
See the License for the specific language governing permissions and
16+
limitations under the License.
17+
18+
-->
19+
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
20+
21+
<modelVersion>4.0.0</modelVersion>
22+
23+
<parent>
24+
<groupId>com.datasqrl.flinkrunner</groupId>
25+
<artifactId>connectors</artifactId>
26+
<version>1.0.0-SNAPSHOT</version>
27+
</parent>
28+
29+
<artifactId>kafka-safe</artifactId>
30+
31+
<dependencies>
32+
<dependency>
33+
<groupId>org.apache.flink</groupId>
34+
<artifactId>flink-core</artifactId>
35+
<version>${flink.version}</version>
36+
<scope>provided</scope>
37+
</dependency>
38+
39+
<dependency>
40+
<groupId>org.apache.flink</groupId>
41+
<artifactId>flink-streaming-java</artifactId>
42+
<version>${flink.version}</version>
43+
<scope>provided</scope>
44+
</dependency>
45+
46+
<dependency>
47+
<groupId>org.apache.flink</groupId>
48+
<artifactId>flink-table-api-java</artifactId>
49+
<version>${flink.version}</version>
50+
<scope>provided</scope>
51+
</dependency>
52+
53+
<dependency>
54+
<groupId>org.apache.flink</groupId>
55+
<artifactId>flink-table-api-java-bridge</artifactId>
56+
<version>${flink.version}</version>
57+
<scope>provided</scope>
58+
</dependency>
59+
60+
<dependency>
61+
<groupId>org.apache.flink</groupId>
62+
<artifactId>flink-connector-kafka</artifactId>
63+
<version>${kafka.version}</version>
64+
<scope>provided</scope>
65+
</dependency>
66+
67+
<dependency>
68+
<groupId>com.google.auto.service</groupId>
69+
<artifactId>auto-service</artifactId>
70+
</dependency>
71+
</dependencies>
72+
</project>

connectors/kafka/src/main/java/com/datasqrl/DeserFailureHandler.java renamed to connectors/kafka-safe/src/main/java/com/datasqrl/DeserFailureHandler.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,9 +27,10 @@
2727
import org.slf4j.Logger;
2828
import org.slf4j.LoggerFactory;
2929

30-
public class DeserFailureHandler {
30+
public class DeserFailureHandler implements Serializable {
3131

3232
private static final Logger LOG = LoggerFactory.getLogger(DeserFailureHandler.class);
33+
private static final long serialVersionUID = 1L;
3334

3435
private final DeserFailureHandlerType handlerType;
3536
private final @Nullable DeserFailureProducer producer;

connectors/kafka/src/main/java/com/datasqrl/DeserFailureHandlerOptions.java renamed to connectors/kafka-safe/src/main/java/com/datasqrl/DeserFailureHandlerOptions.java

File renamed without changes.

connectors/kafka/src/main/java/com/datasqrl/DeserFailureHandlerType.java renamed to connectors/kafka-safe/src/main/java/com/datasqrl/DeserFailureHandlerType.java

File renamed without changes.

connectors/kafka/src/main/java/com/datasqrl/DeserFailureProducer.java renamed to connectors/kafka-safe/src/main/java/com/datasqrl/DeserFailureProducer.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
class DeserFailureProducer implements Serializable {
3131

3232
private static final Logger LOG = LoggerFactory.getLogger(DeserFailureProducer.class);
33+
private static final long serialVersionUID = 1L;
3334

3435
private final String topic;
3536
private final Properties producerProps;

connectors/kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeDynamicKafkaDeserializationSchema.java renamed to connectors/kafka-safe/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeDynamicKafkaDeserializationSchema.java

File renamed without changes.

connectors/kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeKafkaDynamicSource.java renamed to connectors/kafka-safe/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeKafkaDynamicSource.java

File renamed without changes.

connectors/kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeKafkaDynamicTableFactory.java renamed to connectors/kafka-safe/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeKafkaDynamicTableFactory.java

File renamed without changes.

connectors/kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeUpsertKafkaDynamicTableFactory.java renamed to connectors/kafka-safe/src/main/java/org/apache/flink/streaming/connectors/kafka/table/SafeUpsertKafkaDynamicTableFactory.java

File renamed without changes.

connectors/kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/package-info.java renamed to connectors/kafka-safe/src/main/java/org/apache/flink/streaming/connectors/kafka/table/package-info.java

File renamed without changes.

0 commit comments

Comments
 (0)