Class KafkaRecordSerializer
java.lang.Object
com.logicalclocks.hsfs.flink.engine.KafkaRecordSerializer
- All Implemented Interfaces:
Serializable,org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema<org.apache.avro.generic.GenericRecord>
public class KafkaRecordSerializer
extends Object
implements org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema<org.apache.avro.generic.GenericRecord>
- See Also:
-
Nested Class Summary
Nested classes/interfaces inherited from interface org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema
org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema.KafkaSinkContext -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionvoidopen(org.apache.flink.api.common.serialization.SerializationSchema.InitializationContext context, org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema.KafkaSinkContext sinkContext) org.apache.kafka.clients.producer.ProducerRecord<byte[],byte[]> serialize(org.apache.avro.generic.GenericRecord genericRecord, org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema.KafkaSinkContext context, Long timestamp) byte[]serializeKey(org.apache.avro.generic.GenericRecord genericRecord) byte[]serializeValue(org.apache.avro.generic.GenericRecord genericRecord)
-
Constructor Details
-
KafkaRecordSerializer
public KafkaRecordSerializer(StreamFeatureGroup streamFeatureGroup) throws FeatureStoreException, IOException - Throws:
FeatureStoreExceptionIOException
-
-
Method Details
-
open
public void open(org.apache.flink.api.common.serialization.SerializationSchema.InitializationContext context, org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema.KafkaSinkContext sinkContext) - Specified by:
openin interfaceorg.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema<org.apache.avro.generic.GenericRecord>
-
serialize
public org.apache.kafka.clients.producer.ProducerRecord<byte[],byte[]> serialize(org.apache.avro.generic.GenericRecord genericRecord, org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema.KafkaSinkContext context, Long timestamp) - Specified by:
serializein interfaceorg.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema<org.apache.avro.generic.GenericRecord>
-
serializeKey
public byte[] serializeKey(org.apache.avro.generic.GenericRecord genericRecord) -
serializeValue
public byte[] serializeValue(org.apache.avro.generic.GenericRecord genericRecord)
-