public class KafkaDebeziumAvroDeserializationSchema extends Object implements org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord>
CdcSourceRecord
.Constructor and Description |
---|
KafkaDebeziumAvroDeserializationSchema(org.apache.flink.configuration.Configuration cdcSourceConfig) |
Modifier and Type | Method and Description |
---|---|
CdcSourceRecord |
deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> message) |
org.apache.flink.api.common.typeinfo.TypeInformation<CdcSourceRecord> |
getProducedType() |
boolean |
isEndOfStream(CdcSourceRecord nextElement) |
void |
open(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context) |
public KafkaDebeziumAvroDeserializationSchema(org.apache.flink.configuration.Configuration cdcSourceConfig)
public void open(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context) throws Exception
open
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord>
Exception
public CdcSourceRecord deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> message) throws IOException
deserialize
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord>
IOException
public boolean isEndOfStream(CdcSourceRecord nextElement)
isEndOfStream
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord>
public org.apache.flink.api.common.typeinfo.TypeInformation<CdcSourceRecord> getProducedType()
getProducedType
in interface org.apache.flink.api.java.typeutils.ResultTypeQueryable<CdcSourceRecord>
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.