Skip to content

Commit d33ba50

Browse files
authored
fix: Wire SafeUpsertKafkaDynamicTableFactory (#325)
1 parent f0bda49 commit d33ba50

1 file changed

Lines changed: 3 additions & 0 deletions

File tree

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
package org.apache.flink.streaming.connectors.kafka.table;
1717

1818
import com.datasqrl.flinkrunner.connector.kafka.DeserFailureHandler;
19+
import com.google.auto.service.AutoService;
1920
import org.apache.flink.api.common.serialization.DeserializationSchema;
2021
import org.apache.flink.api.common.serialization.SerializationSchema;
2122
import org.apache.flink.api.java.tuple.Tuple2;
@@ -37,6 +38,7 @@
3738
import org.apache.flink.table.factories.DeserializationFormatFactory;
3839
import org.apache.flink.table.factories.DynamicTableSinkFactory;
3940
import org.apache.flink.table.factories.DynamicTableSourceFactory;
41+
import org.apache.flink.table.factories.Factory;
4042
import org.apache.flink.table.factories.FactoryUtil;
4143
import org.apache.flink.table.factories.SerializationFormatFactory;
4244
import org.apache.flink.table.types.DataType;
@@ -82,6 +84,7 @@
8284
import static org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptionsUtil.validateScanBoundedMode;
8385

8486
/** Upsert-Kafka factory. */
87+
@AutoService(Factory.class)
8588
public class SafeUpsertKafkaDynamicTableFactory
8689
implements DynamicTableSourceFactory, DynamicTableSinkFactory {
8790

0 commit comments

Comments
 (0)