diff --git a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaDebeziumAvroDeserializationSchema.java b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaDebeziumAvroDeserializationSchema.java index 6c53dd05d533..dc36151e7bbf 100644 --- a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaDebeziumAvroDeserializationSchema.java +++ b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/action/cdc/kafka/KafkaDebeziumAvroDeserializationSchema.java @@ -41,14 +41,12 @@ public class KafkaDebeziumAvroDeserializationSchema private static final long serialVersionUID = 1L; - private final String topic; private final String schemaRegistryUrl; /** The deserializer to deserialize Debezium Avro data. */ private ConfluentAvroDeserializationSchema avroDeserializer; public KafkaDebeziumAvroDeserializationSchema(Configuration cdcSourceConfig) { - this.topic = KafkaActionUtils.findOneTopic(cdcSourceConfig); this.schemaRegistryUrl = cdcSourceConfig.get(SCHEMA_REGISTRY_URL); } @@ -69,6 +67,7 @@ public void deserialize(ConsumerRecord message, Collector