T - the type to be serializedpublic class AvroSerializationSchema<T> extends Object implements org.apache.flink.api.common.serialization.SerializationSchema<T>
| Modifier | Constructor and Description |
|---|---|
protected |
AvroSerializationSchema(Class<T> recordClazz,
org.apache.avro.Schema schema,
AvroFormatOptions.AvroEncoding encoding)
Creates an Avro deserialization schema.
|
| Modifier and Type | Method and Description |
|---|---|
protected void |
checkAvroInitialized() |
boolean |
equals(Object o) |
static AvroSerializationSchema<org.apache.avro.generic.GenericRecord> |
forGeneric(org.apache.avro.Schema schema)
Creates
AvroSerializationSchema that serializes GenericRecord using provided
schema. |
static AvroSerializationSchema<org.apache.avro.generic.GenericRecord> |
forGeneric(org.apache.avro.Schema schema,
AvroFormatOptions.AvroEncoding encoding)
Creates
AvroSerializationSchema that serializes GenericRecord using provided
schema. |
static <T extends org.apache.avro.specific.SpecificRecord> |
forSpecific(Class<T> tClass)
Creates
AvroSerializationSchema that serializes SpecificRecord using provided
schema. |
static <T extends org.apache.avro.specific.SpecificRecord> |
forSpecific(Class<T> tClass,
AvroFormatOptions.AvroEncoding encoding)
Creates
AvroSerializationSchema that serializes SpecificRecord using provided
schema. |
protected org.apache.avro.generic.GenericDatumWriter<T> |
getDatumWriter() |
protected org.apache.avro.io.Encoder |
getEncoder() |
protected ByteArrayOutputStream |
getOutputStream() |
org.apache.avro.Schema |
getSchema() |
int |
hashCode() |
void |
open(org.apache.flink.api.common.serialization.SerializationSchema.InitializationContext context) |
byte[] |
serialize(T object) |
protected AvroSerializationSchema(Class<T> recordClazz, @Nullable org.apache.avro.Schema schema, AvroFormatOptions.AvroEncoding encoding)
recordClazz - class to serialize. Should be one of: SpecificRecord, GenericRecord.schema - writer Avro schema. Should be provided if recordClazz is GenericRecordpublic static <T extends org.apache.avro.specific.SpecificRecord> AvroSerializationSchema<T> forSpecific(Class<T> tClass)
AvroSerializationSchema that serializes SpecificRecord using provided
schema.tClass - the type to be serializedpublic static <T extends org.apache.avro.specific.SpecificRecord> AvroSerializationSchema<T> forSpecific(Class<T> tClass, AvroFormatOptions.AvroEncoding encoding)
AvroSerializationSchema that serializes SpecificRecord using provided
schema.tClass - the type to be serializedpublic static AvroSerializationSchema<org.apache.avro.generic.GenericRecord> forGeneric(org.apache.avro.Schema schema)
AvroSerializationSchema that serializes GenericRecord using provided
schema.schema - the schema that will be used for serializationpublic static AvroSerializationSchema<org.apache.avro.generic.GenericRecord> forGeneric(org.apache.avro.Schema schema, AvroFormatOptions.AvroEncoding encoding)
AvroSerializationSchema that serializes GenericRecord using provided
schema.schema - the schema that will be used for serializationpublic org.apache.avro.Schema getSchema()
protected org.apache.avro.io.Encoder getEncoder()
protected org.apache.avro.generic.GenericDatumWriter<T> getDatumWriter()
protected ByteArrayOutputStream getOutputStream()
public void open(org.apache.flink.api.common.serialization.SerializationSchema.InitializationContext context)
throws Exception
public byte[] serialize(T object)
serialize in interface org.apache.flink.api.common.serialization.SerializationSchema<T>protected void checkAvroInitialized()
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.