Skip to content
Back to Insights
Data EngineeringBy KE Engineering Team

Flink and Kafka

Flink and KafkaDATA ENGINEERING cover for Flink and KafkaSCHEMA REGISTRYKAFKAFLINKAVRO · SCHEMA IDONE CONTRACT, BOTH SIDESDATA ENGINEERINGFlink and Kafka// SCHEMA REGISTRY · FLINK

First written January 2023, last updated September 2026.

If you are an Apache Kafka developer looking to write stream-processing applications in Flink, the initial setup isn’t so obvious. Apache Flink has its own opinions on consuming and producing to Kafka along with its integration with Confluent’s Schema Registry. Here are the steps, with a working example, to get an Apache Kafka and Apache Flink streaming platform up quickly.

When both the Kafka key and value are part of the streaming pipeline, nested generics come into play and type handling gets tricky.

Resources

  • A full container-based toolkit is available at dev-local. It's based on docker-compose and provides container instances of Apache Kafka, Apache Flink, Kafka Connect, and more.
  • A demonstration of this article is provided in the flink folder of the dev-local-demos project.

Challenges

The key challenges uncovered:

  • Committing Offsets
  • Using Offsets
  • Serialization and Deserialization
  • Confluent’s Schema Registry Integration
  • Java Generics and Type Erasure

Committing Offsets

Flink's KafkaSource commits Kafka consumer offsets when a checkpoint completes, and checkpointing is off by default. So a job without checkpointing never commits on Flink's terms, and when it restarts it consumes from the earliest or latest offset, depending on the reset strategy. The step to remember is turning checkpointing on.

Add a Kafka group.id to your consumer. commit.offsets.on.checkpoint is on by default; the sample sets it explicitly so the intent is visible.

java
KafkaSource<Tuple2<EventKey, EventValue>> source = KafkaSource.<Tuple2<EventKey, EventValue>>builder()    ...    .setProperty("commit.offsets.on.checkpoint", "true")    .setProperty("group.id", "flink-processor")    ...    .build();

Checkpointing must be enabled. It's one line, and it's the line people forget.

java
StreamExecutionEnvironment.getExecutionEnvironment().enableCheckpointing(100L); // demo only; production intervals are usually seconds to minutes

Using Offsets

Committing offsets doesn’t mean they will be used; the KafkaSource instance has to be configured to use them. The committedOffsets initializer of OffsetsInitializer instructs the source instance to use them. Be sure also to define the behavior when offsets don't yet exist by selecting the appropriate OffsetResetStrategy.

java
KafkaSource<Tuple2<EventKey, EventValue>> source = KafkaSource.<Tuple2<EventKey, EventValue>>builder()    ...    .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))    ...    .build();

Serialization and Deserialization

The Kafka Producer and Consumer used by Flink use the byte[] serializer and deserializer and leave the marshaling of data to Flink. If you only need message values, setup is easier. In stream processing, though, a message key is usually necessary, so it's best to understand how that works.

Implement the interfaces KafkaRecordDeserializationSchema and KafkaRecordSerializationSchema, and use them in the KafkaSource and KafkaSink builders, respectively.

The method names within KafkaSource and KafkaSink aren't consistent. For KafkaSource it is setDeserializer and takes the KafkaRecordDeserializationSchema. For KafkaSink it is setRecordSerializer and takes the KafkaRecordSerializationSchema.

java
KafkaSource<Tuple2<EventKey, EventValue>> source = KafkaSource.<Tuple2<EventKey, EventValue>>builder()    ...    .setDeserializer(new MyAvroDeserialization<>(EventKey.class, EventValue.class, "http://schema-registry:8081"))    .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))    ...    .build();
java
KafkaSink<Tuple2<EventKey, EventValue>> sink = KafkaSink.<Tuple2<EventKey, EventValue>>builder()    ...    .setRecordSerializer(new MyAvroSerialization<>(EventKey.class, EventValue.class, topic, "http://schema-registry:8081"))    ...    .build();

Confluent’s Schema Registry Integration

Flink provides its own Confluent schema class, ConfluentRegistryAvroDeserializationSchema, as part of the flink-avro-confluent-registry library. Use this within your MyAvroDeserialization and MyAvroSerialization implementation. The configuration of Schema Registry and subject naming conventions need to be implemented accordingly.

java
ConfluentRegistryAvroDeserializationSchema.forSpecific(pojoClass, schemaRegistryUrl)

For topic-based serialization, the subject would be {topic}-key and {topic}-value.

java
ConfluentRegistryAvroSerializationSchema.forSpecific(pojoClass, subject, schemaRegistryUrl);

Java Generics and Type Erasure

Kafka Streams makes both key and value part of the Processor API and domain-specific language (DSL). This reduces the complexities of using generics. In Flink, the record is a single object. Capturing both Key and Value objects within Flink requires more nuance with generics.

Fortunately, Flink’s Tuple (and its 26 concrete implementations Tuple0, Tuple1, …, Tuple25) make it possible to write serializers and deserializers generically. The core concept is that the serializer and deserializer need to provide type information in a way where the nested types aren't erased. This is done through TypeInformation.

Creating a TypeHint works similarly to TypeReference in the Jackson JSON library. A specific instance allows the code to handle the erased types.

java
TypeInformation<Tuple2<EventKey, EventValue>> typeInformation = TypeInformation.of(    new TypeHint<Tuple2<EventKey, EventValue>>() {    });

With the use of TupleTypeInfo a single instance of KafkaRecordDeserializationSchema and KafkaRecordSerializationSchema can be created for handling marshaling of Avro to SpecificRecord implementations.

java
new TupleTypeInfo<>(TypeInformation.of(keyClass), TypeInformation.of(valueClass));

Putting it all Together

This demonstration is available in the dev-local-demos project, flink.

A simple ETL-based processor with a transformation of an Order to a PurchaseOrder (as shown in the demonstration code), is achieved with a simple functional transformation. In this example, the business logic is the convert method within the map operation.

java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment()    .enableCheckpointing(100L); // demo only; production intervals are usually seconds to minutes KafkaSource<Tuple2<OrderKey, Order>> source = KafkaSource.<Tuple2<OrderKey, Order>>builder()    .setBootstrapServers(bootStrapServers)    .setTopics("input")    .setDeserializer(new AvroDeserialization<>(OrderKey.class, Order.class, schemaRegistryUrl))    .setProperty("commit.offsets.on.checkpoint", "true")    .setProperty("group.id", "FLINK")    .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))    .build(); KafkaSink<Tuple2<OrderKey, PurchaseOrder>> sink = KafkaSink.<Tuple2<OrderKey, PurchaseOrder>>builder()    .setBootstrapServers(bootStrapServers)    .setRecordSerializer(new AvroSerialization<>(OrderKey.class, PurchaseOrder.class, "output", schemaRegistryUrl))    .build(); env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source")    .map(kv -> {        return new Tuple2<>(kv.f0, convert(kv.f1));    }, new TupleTypeInfo<>(TypeInformation.of(OrderKey.class), TypeInformation.of(PurchaseOrder.class)))    .sinkTo(sink); env.execute();

Running your application

Build the application with dependencies bundled and run it with the flink run command. Parallelization is set with the -p option.

shell
flink run -p 4 --detached application-all-dependencies.jar

The one thing to remember

If your job restarts from the beginning, check that checkpointing is on. That's the one people miss.

Working on something like this?

Start a Conversation