What’s actually slowing this PC down?

Pick the symptom - the matching free tool is one click away.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

To consume Kafka messages in Flink, use KafkaSource in a DataStream job or configure the Kafka connector in Table/SQL. Choose the starting offset explicitly and enable checkpointing if the job must recover from failures without losing its coordinated Flink state. The examples and defaults depend on the Flink release, so check the documentation for the version you deploy.

Choose the Flink Kafka API that matches your job

Flink provides separate Kafka integrations for DataStream and Table/SQL applications. They use different APIs and configuration options; do not copy a setting or assume a default from one into the other.

Job type Kafka integration Where to configure it
DataStream KafkaSource Build the source with the DataStream connector API. Starting offsets are selected with an OffsetsInitializer.
Table or SQL Kafka table connector Set the connector and Kafka options in the table definition or corresponding Table API configuration.

Use the documentation for your deployed Flink release to confirm connector compatibility, dependency coordinates, and option names. The official Flink 2.1 DataStream Kafka connector guide and the Kafka Table connector guide are version-specific references; neither establishes which artifact is correct for every Flink and Kafka deployment.

Decide where consumption should begin

Choose a starting position based on the records the job needs to process. Starting at the earliest offset can replay retained history; starting at the latest generally begins with records arriving after the source starts. Resuming from a consumer group’s committed offsets depends on offsets being available and on the configured behavior when they are missing.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Starting position Use it when Important consideration
Committed group offsets Resuming a consumer group’s recorded progress. Set the fallback or reset behavior for partitions without a committed offset; do not assume it is identical across APIs or releases.
Earliest Reading from the oldest offset Kafka still retains. This may replay a substantial backlog, but cannot recover records Kafka has already deleted.
Latest Starting near the current end of the log. Earlier retained records are not part of the initial read.
Timestamp Starting from offsets associated with a chosen time. Confirm the connector’s timestamp semantics and behavior when a partition has no matching offset.
Specific offsets Starting partitions at explicitly selected offsets. Specify and validate offsets for the partitions the job will read.

In DataStream, select an OffsetsInitializer for the intended start position. In Table/SQL, use the documented startup-mode options for that connector release. Avoid relying on an implicit default when replay or missed records would matter.

Choose whether the Kafka read is bounded

A continuously running streaming job usually reads without a predetermined end. A backfill or finite scan may instead need to stop at selected offsets or another supported boundary. The Table connector documents bounded scans and ending modes such as latest, timestamp, group offsets, or specific offsets. Check the options available for your API and release before designing a finite read; do not assume that a Table/SQL stopping option is available in the same form in DataStream.

Use Flink checkpoints for fault-tolerant recovery

For a DataStream source participating in Flink checkpointing, source offsets are included in Flink state. Following a failure, Flink can restore the source position from the recovered checkpoint. The DataStream Kafka connector documentation describes committing offsets to Kafka after completed checkpoints so other consumers and monitoring tools can observe progress; those broker commits are not the source’s recovery mechanism.

Enable checkpointing in the job’s Flink configuration or application code, and verify that checkpoints complete successfully. If checkpointing is disabled, Kafka client auto-commit behavior may be controlled by consumer properties, but broker commits alone do not coordinate source progress with Flink state and should not be treated as equivalent recovery.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Understand what exactly-once means

Exactly-once state updates are not automatically the same as exactly-once delivery through an entire Kafka-to-Flink-to-Kafka pipeline. Flink’s fault-tolerance guarantees documentation says exactly-once updates to user-defined state require the source to participate in snapshotting. End-to-end delivery also depends on the sink’s guarantees and configuration.

For transactional Kafka output, the Table connector describes exactly-once delivery when checkpointing is enabled. A downstream Kafka consumer that must not see uncommitted transactional records should use read_committed isolation. Confirm the relevant source, sink, and consumer settings together rather than inferring an end-to-end guarantee from the fact that the job reads Kafka.

Account for watermarks and idle partitions

In event-time jobs, a Kafka partition that stops producing can hold back downstream watermark progress. In the Flink 2.1 Kafka connector documentation, source parallelism greater than the number of Kafka partitions does not by itself make unused source readers idle. Configure idleness in the watermark strategy when appropriate so idle inputs do not hold back watermarks, and confirm the exact setting for the API and release in use.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Monitor the running source

Track source progress and consumer lag to distinguish a healthy stream from a stalled or falling-behind job. Interpret Kafka’s committed offsets as visible consumer progress, not as proof that Flink has a recoverable checkpoint. Check metric names and availability against the connector release deployed with the job.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.