public class KafkaLogDeserializationSchema extends Object implements org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<org.apache.flink.table.data.RowData>
KafkaDeserializationSchema
for the table with primary key in log store.Constructor and Description |
---|
KafkaLogDeserializationSchema(org.apache.flink.table.types.DataType physicalType,
int[] primaryKey,
org.apache.flink.api.common.serialization.DeserializationSchema<org.apache.flink.table.data.RowData> primaryKeyDeserializer,
org.apache.flink.api.common.serialization.DeserializationSchema<org.apache.flink.table.data.RowData> valueDeserializer,
int[][] projectFields) |
Modifier and Type | Method and Description |
---|---|
org.apache.flink.table.data.RowData |
deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> record) |
void |
deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> record,
org.apache.flink.util.Collector<org.apache.flink.table.data.RowData> underCollector) |
org.apache.flink.api.common.typeinfo.TypeInformation<org.apache.flink.table.data.RowData> |
getProducedType() |
boolean |
isEndOfStream(org.apache.flink.table.data.RowData nextElement) |
void |
open(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context) |
public KafkaLogDeserializationSchema(org.apache.flink.table.types.DataType physicalType, int[] primaryKey, @Nullable org.apache.flink.api.common.serialization.DeserializationSchema<org.apache.flink.table.data.RowData> primaryKeyDeserializer, org.apache.flink.api.common.serialization.DeserializationSchema<org.apache.flink.table.data.RowData> valueDeserializer, @Nullable int[][] projectFields)
public void open(org.apache.flink.api.common.serialization.DeserializationSchema.InitializationContext context) throws Exception
open
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<org.apache.flink.table.data.RowData>
Exception
public boolean isEndOfStream(org.apache.flink.table.data.RowData nextElement)
isEndOfStream
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<org.apache.flink.table.data.RowData>
public org.apache.flink.table.data.RowData deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> record)
deserialize
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<org.apache.flink.table.data.RowData>
public void deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> record, org.apache.flink.util.Collector<org.apache.flink.table.data.RowData> underCollector) throws Exception
deserialize
in interface org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<org.apache.flink.table.data.RowData>
Exception
public org.apache.flink.api.common.typeinfo.TypeInformation<org.apache.flink.table.data.RowData> getProducedType()
getProducedType
in interface org.apache.flink.api.java.typeutils.ResultTypeQueryable<org.apache.flink.table.data.RowData>
Copyright © 2023–2025 The Apache Software Foundation. All rights reserved.