Skip to content

Spark Structured Streaming Kafka offsets and lag: where Spark keeps its position

Spark
Chad Harris·October 6, 2026·9 min read

A team runs a Spark Structured Streaming job that reads a Kafka topic. The Spark UI says the query is active, but the consumer group dashboard in the Kafka tooling shows no lag for it, or a group nobody recognises with a name like spark-kafka-source- followed by a long identifier. The natural conclusion is that the job is not reading, or that the monitoring is broken. In most cases neither is true. Spark keeps its read position somewhere else.

Spark Structured Streaming tracks Kafka offsets in its own checkpoint and does not commit them to Kafka. The Apache Spark Kafka integration guide says it directly: for enable.auto.commit, “Kafka source doesn’t commit any offset”, and Structured Streaming “manages which offsets are consumed internally, rather than rely on the kafka Consumer to do it”. A tool that computes lag from the offsets a consumer group has committed has nothing to compare against the end of the topic, so lag for a Spark job has to come from Spark.

Flink has a similar split between its checkpoints and Kafka’s committed offsets, covered in Flink Kafka source offsets and lag. For the Kafka-side basics, see Kafka offsets in the complete Kafka guide.

This page covers where the position lives, why a consumer group view misleads, how to read lag from Spark itself, and how to restart a query without losing or repeating data.

Where Spark keeps its position

A streaming query started with a checkpointLocation saves its progress there. The Structured Streaming programming guide describes this as saving “all the progress information (i.e. range of offsets processed in each trigger)” and the running aggregates to a path on an HDFS compatible file system. For a Kafka source, the progress information is a per-partition offset range for each micro-batch.

Two consequences follow from that.

  • The checkpoint is the only record of the read position. Delete it or point the query at a new location and the query starts again from startingOffsets, which defaults to latest for a streaming query. Rows that arrived while the job was down are skipped without an error.
  • startingOffsets is ignored on a restart. The guide notes that for a streaming query it “only applies when a new query is started, and that resuming will always pick up from where the query left off”. Changing the option on an existing checkpoint does not rewind the job.

Why a consumer group view shows no lag

By default each query generates its own consumer group id. The guide explains that this “ensures that each Kafka source has its own consumer group that does not face interference from any other consumer”. The prefix is spark-kafka-source unless you set the groupIdPrefix option. So every query, and every restart that creates a new query, can appear as a new group name.

You can force a fixed name with the kafka.group.id option, which the guide says to “use with caution”. Concurrently running queries with the same group id “are likely interfere with each other causing each query to read only part of the data”. The usual reason to set it is group-based authorization, where the brokers only allow a specific group name. Setting it does not make Spark commit offsets, so it gives monitoring a stable name and nothing else.

That is why the two views disagree. The Kafka side reports what consumers committed. Spark reports what the query processed, and the two are separate records.

Measure lag from Spark

Every streaming query reports progress. In the programming guide, query.lastProgress and query.recentProgress return the most recent updates, and each source entry carries an endOffset per topic and partition together with numInputRows, inputRowsPerSecond and processedRowsPerSecond. A query that processes fewer rows per second than arrive is falling behind.

For the Kafka source specifically, Spark’s source code adds a metrics block to that progress report. The KafkaMicroBatchStream source file builds minOffsetsBehindLatest, maxOffsetsBehindLatest and avgOffsetsBehindLatest, which are the number of offsets the query is behind the latest offset in each partition. These are the closest thing Spark has to consumer lag, and maxOffsetsBehindLatest is the one to alert on, because a single stalled partition hides inside the average.

To push these into your monitoring instead of polling, attach a StreamingQueryListener with sparkSession.streams.addListener(). The programming guide describes callbacks when a query starts, stops and makes progress, and the progress event carries the same source metrics.

Measure lag from the Kafka side

The same number can be computed from outside Spark. Take the endOffset of each partition from the last progress report, ask the cluster for the latest offset of each partition of the topic, and subtract. This works for any Spark deployment because it needs only the progress report and a read of the topic. It also catches a case Spark metrics can hide, a query that has stopped reporting progress at all.

When lag turns into data loss

A stalled job is not only late. Kafka deletes old records according to retention, so if the position falls behind the oldest retained offset the data is gone. The Kafka integration guide covers this under failOnDataLoss, which defaults to true and fails the query “when it’s possible that data is lost (e.g., topics are deleted, or offsets are out of range)”. The guide warns this “may be a false alarm” and can be switched off, but switching it off turns a loud failure into silent gaps. In extreme cases, the guide says, the input rows of a batch “might be gradually reduced until zero” when the job cannot keep up with retention.

Compare how long the job takes to recover against the topic’s retention before the incident, not during it. The same retention question applies to a Flink job restarting from its checkpoint.

Restart safely

  • Resuming after a failure. Restart with the same checkpointLocation. Spark resumes from the saved offsets and startingOffsets has no effect.
  • Reprocessing from an earlier point on purpose. Start a new query with a new checkpoint location and set startingOffsets or startingTimestamp to the point you want. Spark does this by creating a new query, which also creates a new consumer group id.
  • Never reuse one kafka.group.id across queries. Two queries sharing a group id can each read only part of the data.

Where Kpow fits

This page is published by Factor House, which makes Kpow, its Kafka management product, so weigh this section with that in mind. Factor House does not sell a Spark product. Kpow does not read Spark checkpoints or Spark progress reports, and it cannot show where a Structured Streaming query is in a topic.

What Kpow does cover is the Kafka side of the same pipeline. Its group management documentation describes inspecting consumer groups and changing their offsets, which applies to consumers that do commit offsets, including the other jobs reading the same topic. The offset management tools comparison and the consumer lag monitoring tools comparison rank the tools for that side of the pipeline. For a consumer that does commit, Kpow changes offsets only when the group is in the EMPTY state, and that is not a way to rewind a Spark query, whose position is in the checkpoint.

If you want to see a Spark Structured Streaming job reading Kafka beside Kpow, Factor House publishes Lab 10 of Factor House Local. It builds a PySpark job that reads Kafka order records and writes them to an Iceberg table, with Kpow running against the Kafka environment.

FAQ

Does Spark Structured Streaming commit offsets to Kafka?

No. The Spark documentation states that the Kafka source does not commit any offset. Progress is saved in the checkpoint location, so lag has to be read from the query progress or computed against the topic.

Why does my Spark job have a different consumer group name after every restart?

Each query generates its own group id from the groupIdPrefix, which is spark-kafka-source by default. A restart that starts a new query can produce a new name. Set kafka.group.id only if brokers require a specific group, and never share it between concurrent queries.

How do I measure consumer lag for a Spark streaming job?

Read maxOffsetsBehindLatest from the Kafka source metrics in the query progress, or subtract the endOffset in the progress report from the latest offset of each partition. Alert on the maximum across partitions.

Can I change startingOffsets to replay data?

Not on an existing checkpoint. The setting only applies when a new query starts. To replay, start a new query with a new checkpoint location.

Related reading