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

Connectors: Sources and Sinks

An overview of Flink's connector ecosystem for reading from and writing to external systems like Kafka, files, and databases.

DataStream APIIntermediate9 min readJul 10, 2026
Analogies

Modern Flink jobs read data using the unified Source interface introduced by FLIP-27, which splits the work of ingestion into a SplitEnumerator running on the JobManager that discovers and assigns work units (splits), and one or more SourceReaders running on TaskManagers that actually read records from those splits, such as individual Kafka partitions or file shards. This design unifies bounded and unbounded reading under one API and lets connectors like KafkaSource handle dynamic partition discovery, checkpointed offsets, and backpressure-aware reading without custom code in the job itself.

🏏

Cricket analogy: The SplitEnumerator is like a tournament scheduler assigning specific matches to specific stadiums, while each stadium's ground staff (SourceReader) is responsible for actually running the fixtures assigned to it.

Sinks: Getting Data Out

The counterpart Sink interface defines how records leave a Flink job, typically through a SinkWriter that batches and writes records to an external system, optionally coordinated by a Committer and GlobalCommitter for connectors that support two-phase commit, such as KafkaSink configured with DeliveryGuarantee.EXACTLY_ONCE. In that mode, writes are staged inside a transaction tied to a checkpoint, and only committed once Flink confirms the checkpoint has completed successfully across the whole job, which is what allows exactly-once output semantics even after a failure and restart.

🏏

Cricket analogy: A two-phase commit sink is like a third umpire's decision that stays 'pending' on the big screen until confirmed by the review process, only becoming official once every check has passed — not the instant the on-field umpire raises a finger.

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

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

env.fromSource(source, WatermarkStrategy.noWatermarks(), "orders-source")
   .map(order -> enrich(order))
   .sinkTo(sink);

env.execute("Kafka To Kafka Enrichment");

Delivery Guarantees

The end-to-end delivery guarantee of a Flink pipeline depends on the combination of the source's ability to replay records from a checkpointed offset, Flink's own checkpointing to make operator state consistent, and the sink's own guarantee, which ranges from at-most-once (no retry on failure), through at-least-once (retries can produce duplicates), to exactly-once (achieved via idempotent writes or two-phase commit transactions). A pipeline is only as strong as its weakest link, so pairing an exactly-once-capable source like Kafka with a non-transactional, non-idempotent sink still only yields at-least-once semantics overall.

🏏

Cricket analogy: It's like a run being counted correctly only if both the batsmen actually complete the run and the umpire confirms it — a fast start with a run-out at the far end still means the run doesn't count, no matter how good the start was.

Two-phase commit sinks rely on Flink's checkpoint barriers to know when it's safe to commit a transaction: the sink pre-commits (stages) writes as part of each checkpoint, and only actually commits them once the checkpoint has been acknowledged as complete by every operator in the job, tying external write durability directly to Flink's own checkpointing protocol.

Custom Connectors

When no existing connector fits a bespoke system, Flink allows implementing a custom Source by providing a SplitEnumerator that discovers and assigns splits (e.g. partitions of a proprietary message queue) and a SourceReader that reads records from an assigned split and reports its progress for checkpointing, or a custom Sink by implementing a SinkWriter and, if exactly-once is required, a Committer that finalizes staged writes only after a successful checkpoint. This is more involved than using an existing connector, but it lets any bespoke external system participate correctly in Flink's checkpointing and fault-tolerance model.

🏏

Cricket analogy: Writing a custom connector is like a franchise building its own scouting network from scratch in a country with no existing cricket board infrastructure, rather than simply plugging into the well-established BCCI pipeline.

Under at-least-once delivery, retries after a failure can cause the same record to be written more than once. If your sink is not naturally idempotent — for example, a plain INSERT into a relational table without a unique constraint on a business key — duplicate writes will silently corrupt downstream data even though the pipeline appears healthy.

  • Flink's modern Source interface (FLIP-27) splits ingestion into a SplitEnumerator and per-subtask SourceReaders.
  • Sinks write via a SinkWriter, optionally coordinated by a Committer for two-phase commit exactly-once writes.
  • End-to-end delivery guarantees depend on the weakest link among source replay, checkpointing, and sink behavior.
  • Two-phase commit sinks tie transaction commits directly to Flink's checkpoint completion.
  • Custom Sources and Sinks let bespoke external systems participate correctly in checkpointing.
  • At-least-once delivery can produce duplicates unless the sink is idempotent.
  • Choosing the right delivery guarantee requires understanding both source and sink capabilities together.

Practice what you learned

Was this page helpful?

Topics covered

#Programming#ApacheFlinkStudyNotes#ConnectorsSourcesAndSinks#Connectors#Sources#Sinks#Data#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