Package co.decodable.sdk.pipeline.serde
Class DecodableRecordDeserializationSchema<T extends DecodableKeyedStreamRecord<?,?>>
java.lang.Object
co.decodable.sdk.pipeline.serde.DecodableRecordDeserializationSchema<T>
- All Implemented Interfaces:
Serializable,org.apache.flink.api.java.typeutils.ResultTypeQueryable<T>,org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema<T>
public final class DecodableRecordDeserializationSchema<T extends DecodableKeyedStreamRecord<?,?>>
extends Object
implements org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema<T>
- See Also:
-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionvoiddeserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[], byte[]> record, org.apache.flink.util.Collector<T> out) org.apache.flink.api.common.typeinfo.TypeInformation<T> voidopen(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context)
-
Constructor Details
-
DecodableRecordDeserializationSchema
-
-
Method Details
-
deserialize
public void deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[], byte[]> record, org.apache.flink.util.Collector<T> out) throws IOException- Specified by:
deserializein interfaceorg.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema<T extends DecodableKeyedStreamRecord<?,?>> - Throws:
IOException
-
getProducedType
- Specified by:
getProducedTypein interfaceorg.apache.flink.api.java.typeutils.ResultTypeQueryable<T extends DecodableKeyedStreamRecord<?,?>>
-
open
public void open(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context) throws Exception - Specified by:
openin interfaceorg.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema<T extends DecodableKeyedStreamRecord<?,?>> - Throws:
Exception
-