Avro

Updated at:

The Avro format reads and writes data based on an Avro schema automatically derived from the table schema. Supported connectors include the Kafka connector, Upsert Kafka connector, and Object Storage Service (OSS) connector.

Example

The following DDL creates a table that reads from and writes to a Kafka topic in Avro format.

CREATE TABLE user_behavior (
  user_id     BIGINT,
  item_id     BIGINT,
  category_id BIGINT,
  behavior    STRING,
  ts          TIMESTAMP(3)
) WITH (
  'connector'                    = 'kafka',
  'topic'                        = 'user_behavior',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id'          = 'testGroup',
  'format'                       = 'avro'
);

Parameters

ParameterRequiredDefaultData typeDescription
formatYesNoneSTRINGThe format to use. Set to avro.
avro.codecNoNoneSTRINGThe compression codec for Avro files. Takes effect only when the file system connector is used. Valid values: snappy (default when codec is set), null, deflate, bzip2, xz.

Data type mappings

The Avro schema is always derived from the Flink SQL table schema. The following table lists the mappings from Flink SQL types to Avro types.

Flink SQL typeAvro typeAvro logical typeNotes
CHAR / VARCHAR / STRINGstring
BOOLEANboolean
BINARY / VARBINARYbytes
DECIMALbytesdecimalEncodes precision and scale.
TINYINTint
SMALLINTint
INTint
BIGINTlong
FLOATfloat
DOUBLEdouble
DATEintdateValue represents days since the Unix epoch.
TIMEinttime-millisValue represents milliseconds since midnight.
TIMESTAMPlongtimestamp-millisValue represents milliseconds since the Unix epoch.
ARRAYarray
MAP (key must be STRING, CHAR, or VARCHAR)map
MULTISET (element must be STRING, CHAR, or VARCHAR)map
ROWrecord

Flink can read and write Nullable data types. A Nullable type maps to Avro union(something, null), where something is the Avro type corresponding to the Flink SQL type.

For more information about Avro types, see the Avro specification.