100% Free Forever
AI-Powered Learning
Industry Expert Content
Certificates & Badges
Learn At Your Own Pace
Programming

Flink with Kafka

Understand how to build reliable, exactly-once streaming pipelines by connecting Apache Flink to Apache Kafka as both source and sink.

Table API & ProductionIntermediate10 min readJul 10, 2026
Analogies

Kafka provides Flink with a durable, replayable, partitioned log, which matches Flink's own execution model closely: Kafka partitions map naturally onto parallel source subtasks, Kafka's offset-based replay lets Flink restart from a checkpoint without data loss, and Kafka's retention window gives Flink time to recover from failures without the source discarding data. This alignment is why nearly every production Flink deployment uses Kafka (or a Kafka-compatible log like Redpanda or Amazon MSK) as its primary ingestion layer.

🏏

Cricket analogy: It's like a stadium's ball-by-ball commentary archive — because every delivery is logged with a sequence number, a broadcaster who loses connection can resume exactly from ball 247 rather than guessing where they left off, just as Flink resumes from a Kafka offset.

Configuring the Kafka Source and Sink

The modern KafkaSource builder (replacing the deprecated FlinkKafkaConsumer) lets you set bootstrap servers, topic(s) or a topic pattern, a deserializer, and a starting offset strategy such as OffsetsInitializer.earliest(), .latest(), or .committedOffsets(). On the write side, KafkaSink supports three delivery guarantees — NONE, AT_LEAST_ONCE, and EXACTLY_ONCE — with exactly-once implemented via Kafka's transactional producer API, meaning Flink writes records inside a Kafka transaction that only commits when the corresponding checkpoint completes.

🏏

Cricket analogy: Choosing a starting offset is like a substitute umpire joining mid-match and deciding whether to review only from the current over (latest), from ball one (earliest), or from the last confirmed decision (committed) before making calls.

java
KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setTopics("orders")
    .setGroupId("flink-orders-consumer")
    .setStartingOffsets(OffsetsInitializer.committedOffsets())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

KafkaSink<String> sink = KafkaSink.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setRecordSerializer(KafkaRecordSerializationSchema.builder()
        .setTopic("orders-enriched")
        .setValueSerializationSchema(new SimpleStringSchema())
        .build())
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("orders-enricher")
    .build();

DataStream<String> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "kafka-orders");
stream.sinkTo(sink);

End-to-End Exactly-Once Semantics

Achieving true end-to-end exactly-once requires all three legs to cooperate: Kafka source offsets are stored as part of Flink's checkpoint state (not committed independently to Kafka), Flink's internal state is snapshotted via the checkpoint barrier mechanism, and the KafkaSink pre-commits its transaction during the checkpoint and only finalizes the commit once the checkpoint is confirmed complete — this two-phase-commit protocol is why EXACTLY_ONCE delivery requires checkpointing to be enabled and requires downstream Kafka consumers to read with isolation.level=read_committed to avoid seeing uncommitted, potentially-rolled-back data.

🏏

Cricket analogy: It's like a third umpire's DRS decision — the on-field call is provisional (pre-committed) until the TV replay confirms it, and only then does the scoreboard (Kafka log) officially update; viewers with the replay feed on 'confirmed only' mode are like read_committed consumers.

Set isolation.level=read_committed on any downstream Kafka consumer reading a topic written by an EXACTLY_ONCE KafkaSink; otherwise consumers may read uncommitted records from transactions that are later aborted during a Flink failure recovery.

EXACTLY_ONCE Kafka sinks hold open transactions between checkpoints. If your checkpoint interval is long and a consumer's transaction.timeout.ms is shorter than that interval, Kafka will abort the transaction before Flink commits it, silently dropping data — always align these two settings.

Schema Evolution with Kafka Topics

Because Kafka stores raw bytes, schema management is external to the log itself — most production setups pair Flink with a Schema Registry (Confluent or Apicurio) and use Avro or Protobuf formats so that both producers and Flink's deserializer validate against a shared, versioned schema, allowing backward-compatible field additions without breaking already-running consumers.

🏏

Cricket analogy: It's like the ICC maintaining a central rulebook that all member boards reference — a new law addition (like the concussion substitute rule) is versioned and rolled out so old matches' scorecards remain valid under the rules they were played under.

  • Kafka's partitioned, replayable log model aligns naturally with Flink's parallel, checkpoint-based recovery model.
  • KafkaSource replaces the deprecated FlinkKafkaConsumer and supports configurable starting offsets (earliest, latest, committedOffsets).
  • KafkaSink supports NONE, AT_LEAST_ONCE, and EXACTLY_ONCE delivery guarantees.
  • EXACTLY_ONCE relies on Kafka transactions coordinated with Flink's checkpoint barriers via two-phase commit.
  • Downstream consumers must use isolation.level=read_committed to avoid seeing data from aborted transactions.
  • Checkpoint interval and Kafka's transaction.timeout.ms must be aligned to avoid silent data loss from transaction expiry.
  • Schema Registry with Avro/Protobuf enables backward-compatible schema evolution between producers and Flink consumers.

Practice what you learned

Was this page helpful?

Topics covered

#Programming#ApacheFlinkStudyNotes#FlinkWithKafka#Flink#Kafka#Most#Common#StudyNotes#SkillVeris#ExamPrep

Frequently Asked Questions

21 categories · pick one to explore

Where can I get free study notes for programming and tech subjects?
SkillVeris offers completely free study notes covering programming and tech subjects, with no signup fees or paywalls. The notes are structured by course and topic, written for quick understanding, and enriched with the Learn Through Hobbies analogy method, so you can revise concepts through cricket, music, gaming, cooking and more.
Are SkillVeris study notes good for exam revision?
Yes, the study notes are designed for efficient revision: each topic answers its heading immediately, keeps explanations concise, and links to related glossary terms and cheat sheets. Students preparing for university exams or certification tests use them as quick revision notes because they distil concepts without the padding of full textbooks.
What subjects do the free study notes cover?
The study notes span the platform's main domains, including AI and machine learning, Python and programming, web development, DevOps, cloud, security and databases. Coverage mirrors the 37 live courses, so notes exist for the topics you are actually studying, and new note sets are added as courses launch.
How are SkillVeris study notes different from regular textbooks?
The notes are answer-first, concise and free, whereas textbooks are long and often expensive. Each section explains one concept directly, then reinforces it through selectable hobby analogies like cricket or cooking. Notes also cross-link to the glossary, blog and cheat sheets, letting you jump to related material instantly instead of flipping pages.
Can I use the developer study material without creating an account?
The study notes are free to access, and SkillVeris does not charge anything for its developer study material at any point. Browsing notes is straightforward from the Study Notes section, and if you want progress tracking, certificates and AI Mentor conversations tied to your learning, a free account unlocks those extras.
Do the study notes explain concepts with analogies?
Yes, this is a signature SkillVeris feature. Study notes use the Learn Through Hobbies method, explaining technical concepts through analogies from twelve domains including cricket, music, gaming, photography, travel, movies, fitness, chess, cooking, finance, business and sports. You can switch the analogy domain instantly to whichever hobby makes the concept click.
Are the revision notes suitable for last-minute exam preparation?
Yes, revision notes on SkillVeris work well for last-minute preparation because every section states the answer in its first sentences, so skimming is genuinely effective. Pair them with the relevant cheat sheet for formulas and syntax, and use the glossary for any unfamiliar term you meet while cramming.
Is there free study material for AI and machine learning?
Yes, SkillVeris provides free study notes across its AI and ML catalogue, covering Python for AI, deep learning frameworks like PyTorch and TensorFlow, Hugging Face Transformers, Large Language Models, RAG, AI agents and MLOps. All of it is free, making it a strong resource for Indian students and global learners alike.
Can beginners understand the study notes, or are they for experts?
Beginners can absolutely use them. The notes are written in plain language, define terms as they appear, and lean on hobby analogies to make abstract ideas concrete. Difficulty scales with the underlying course level, so beginner-course notes stay gentle while advanced-course notes go deeper, and the glossary supports you throughout.
How do study notes connect with SkillVeris courses?
Study notes are organised by course and topic, so they map directly to the structured courses and their 24–40-lesson curriculum. Many learners study a lesson first, then use the matching notes for revision before module assessments and the final exam, where 80 percent is required to pass and earn the certificate.
Are there study notes for Python specifically?
Yes, Python is well covered through notes tied to the Python-focused courses, including Python for AI and ML. Topics span fundamentals through applied machine learning usage. You can reinforce the notes with Python practice in Code Lab, which runs code in your browser with no installation required.
Do the study notes include code examples?
Yes, study notes include code examples wherever a concept is best shown in code, alongside explanations, key points and analogies. Reading a snippet in the notes and then reproducing it yourself in Code Lab is an effective loop, since Code Lab lets you run code in the browser across six languages.
How often is new study material added to SkillVeris?
Study material grows alongside the course catalogue. Whenever new courses join the platform's 37 live courses, matching study notes, glossary entries and cheat sheets are added so the resources stay in sync. Existing notes are also refined over time, so it is worth revisiting topics you studied earlier.
Can I use SkillVeris notes to prepare for technical interviews?
Yes, the notes make excellent interview revision because they compress each concept into direct, answer-first explanations, which mirrors how you should answer interview questions. Combine them with the SkillVeris interview questions feature, which includes readiness scoring, to test whether your revision has actually made you interview-ready.
Are the study notes mobile-friendly for studying on the go?
Yes, the study notes are built to load fast and read comfortably on mobile devices, so you can revise during a commute or between classes. Sections are short and answer-first, which suits small screens, and analogy switching works on mobile too, letting you study anywhere without carrying books.
What is the difference between study notes and cheat sheets?
Study notes explain concepts in depth with context, examples and analogies, making them ideal for learning and revision. Cheat sheets are compact quick-reference summaries of syntax, commands and key facts, ideal once you already understand a topic. Most learners study the notes first, then keep the cheat sheet handy while coding.
Do study notes help if I am stuck on a course lesson?
Yes, reading the matching study notes often clarifies a lesson because the same concept is explained from a different angle, frequently with a different analogy. If you are still stuck, ask the AI Mentor, which answers 24/7 at Quick, Detailed or Deep-dive depth until the idea genuinely makes sense.
Is there free study material for DevOps and cloud topics?
Yes, SkillVeris carries free study notes for DevOps and cloud topics as part of its coverage across 37 live courses. The material suits learners following the DevOps Engineer or Cloud Engineer paths, and it links to related glossary terms and cheat sheets so you can revise the whole toolchain in one place.
Can school or college students in India use these notes for projects?
Yes, students across India and worldwide use SkillVeris notes for coursework, projects and exam preparation, and everything is free, which matters for student budgets. The notes explain concepts clearly enough to cite in project reports, and Code Lab lets you prototype the project code directly in your browser.
How should I combine study notes with other SkillVeris resources?
A proven loop: learn from a course lesson, revise with the matching study notes, look up unfamiliar terms in the glossary, keep the cheat sheet open while practising in Code Lab, and quiz yourself with interview questions. The AI Mentor fills any remaining gaps 24/7, at whatever depth you need.

What Learners Say

Real journeys from the SkillVeris community — swipe for more.

SkillVeris taught me Python through Cricket. Now I’m building real projects and feeling confident!
Arjun S. · B.Tech Student
The best platform for hobby-based learning. Concepts finally stick.
Priya R. · Data Analyst
I went from zero coding to a portfolio of projects — all by learning through my love for gaming. Landed my first internship!
Kabir M. · CS Undergraduate
Trending Topics50 popular tags — tap to explore
Trending CoursesAll 37 free courses — tap to browse