Flink and Kafka
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.
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.
StreamExecutionEnvironment.getExecutionEnvironment().enableCheckpointing(100L); // demo only; production intervals are usually seconds to minutesUsing 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.
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.
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();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.
ConfluentRegistryAvroDeserializationSchema.forSpecific(pojoClass, schemaRegistryUrl)For topic-based serialization, the subject would be {topic}-key and {topic}-value.
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.
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.
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.
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.
flink run -p 4 --detached application-all-dependencies.jarThe 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