Structured Streaming's Kafka source looks like a single call, spark.readStream.format("kafka"), but it is a small distributed system of its own. On the driver it decides, every micro-batch, exactly which offsets of which partitions to read. On the executors it keeps pools of Kafka consumers and prefetched records alive between batches. Its options decide where a new query starts, how fast it may catch up, how many tasks read a partition, and whether a gap in the log is fatal.

This article treats the source as a component. It walks through how a micro-batch is planned and executed, documents the options as listed for Spark 4.2, works through a catch-up sizing example and lists the failure modes you will meet in production. Checkpoint layout and end-to-end delivery semantics are covered in Structured Streaming + Kafka architecture; this page links there rather than repeating it.

How a micro-batch reads Kafka

Each micro-batch follows the same cycle. The driver asks Kafka for the latest offset of every assigned partition (since Spark 3.1 through an AdminClient, which only needs the Describe ACL on topics; spark.sql.streaming.kafka.useDeprecatedOffsetFetching=true restores the old consumer-based fetch). It applies admission control to cap how far the batch may advance, then writes the chosen end offsets to the checkpoint's offsets/ log before any data is read. That write-ahead step is what makes a restart deterministic: a crashed batch is replayed with exactly the same ranges.

The driver then creates tasks, normally one per Kafka partition, each covering [start, end) for that partition. On an executor the task borrows a consumer from a pool keyed by group, topic and partition, assigns it the partition (the source never uses Kafka's group rebalancing), seeks to the start offset and polls until it reaches the end. When the sink finishes, the driver writes the batch to commits/. The source never commits offsets to Kafka itself, so consumer-group lag tools show nothing useful unless you export progress yourself.

Driver: offset planningAdminClient: latest offsetsrangeoffsets/ logbatch N: start..endplan tasksOne task per offset rangeminPartitions may split a partitionExecutor 1consumer pool + fetched dataExecutor 2assign, seek, pollExecutor 3rows: key, value, offset...Sink writes, then commits/Kafka group offsets untouchedKafka brokerspartitions 0..P-1
The driver fixes each batch's offset ranges in the checkpoint before execution; executors read those ranges with pooled consumers; the commit log marks completion. Kafka-side group offsets are never written.

Where a query starts

Where a query begins is set by three options with a fixed precedence: startingTimestamp (one timestamp for all partitions) beats startingOffsetsByTimestamp (a JSON map of per-partition timestamps), which beats startingOffsets. startingOffsets accepts earliest, latest (the streaming default) or a JSON map where -2 means earliest and -1 means latest. If no offset exists at or after a requested timestamp, startingOffsetsByTimestampStrategy decides: error (default) fails the query, latest starts at the end.

The single most important fact: these options apply only when a query starts with an empty checkpoint. On restart the source resumes from the checkpoint and ignores them. Changing startingOffsets and restarting does nothing; to reprocess, start a new query with a new checkpoint location. Partitions that appear while a query is running, from a topic expansion or a new topic matching subscribePattern, are read from the earliest offset.

from pyspark.sql import functions as F, types as T

schema = T.StructType([
    T.StructField("order_id", T.StringType()),
    T.StructField("amount", T.DoubleType()),
    T.StructField("ts", T.TimestampType()),
])

raw = (spark.readStream.format("kafka")
       .option("kafka.bootstrap.servers", "b1:9093,b2:9093")
       .option("subscribe", "orders")
       .option("startingOffsets", "earliest")         # first start only
       .option("maxOffsetsPerTrigger", 1_200_000)     # cap per micro-batch
       .option("failOnDataLoss", "true")
       .option("includeHeaders", "true")
       .option("kafka.security.protocol", "SASL_SSL")
       .option("kafka.isolation.level", "read_committed")
       .load())

orders = (raw
          .select(F.col("partition"), F.col("offset"),
                  F.from_json(F.col("value").cast("string"), schema).alias("o"))
          .select("partition", "offset", "o.*"))

query = (orders.writeStream
         .option("checkpointLocation", "s3://lake/chk/orders_v1")
         .trigger(processingTime="30 seconds")
         .toTable("bronze.orders"))

Keys and values are always bytes; the source fixes the deserializers. Cast to string and parse, or use from_avro / from_protobuf from the matching Spark packages. Keep partition and offset in the output: they are the cheapest possible lineage and dedup key.

The options that matter

OptionDefaultWhat it controls
subscribe / subscribePattern / assignone requiredTopic list, Java regex, or explicit JSON partition map
failOnDataLosstrueFail when offsets are out of range or a topic vanished
maxOffsetsPerTriggernoneUpper bound on offsets per batch, split across partitions
minOffsetsPerTriggernoneWait for at least this many new offsets before a batch
maxTriggerDelay15mRun a batch anyway after this delay (with min offsets)
minPartitionsnoneSplit large partitions into more Spark tasks
kafkaConsumer.pollTimeoutMs120000Executor poll timeout
fetchOffset.numRetries / retryIntervalMs3 / 10Driver offset-fetch retries
groupIdPrefixspark-kafka-sourcePrefix of the generated consumer group id
includeHeadersfalseAdd a headers array column

Kafka client settings are passed with a kafka. prefix, so security, SASL, TLS and kafka.isolation.level work as on any consumer. Several are refused outright because they would fight the source: group.id (use groupIdPrefix, or kafka.group.id if a broker ACL demands a fixed group, accepting that two queries sharing it interfere), auto.offset.reset (use startingOffsets), enable.auto.commit, the key and value deserializers, and interceptor.classes. Recent documentation also lists maxRecordsPerPartition; check that your Spark version supports it before relying on it.

Admission control and triggers

Without a cap, the first batch after a long outage tries to read the whole backlog in one go, which can mean hours of work in one batch, executor memory pressure and a huge state store commit. maxOffsetsPerTrigger sets a ceiling on the total offsets a batch may cover. The source divides it across partitions in proportion to each partition's backlog, so a lagging partition gets a larger share.

The opposite problem, many tiny batches each paying fixed costs, is solved by minOffsetsPerTrigger: a trigger fires only when at least that many new offsets are available, or when maxTriggerDelay has passed since the last batch, whichever comes first. This is useful for sinks that dislike small files.

For scheduled jobs, Trigger.AvailableNow processes everything available at start-up and then stops, while still honouring maxOffsetsPerTrigger to break the work into multiple batches. It replaces the deprecated Trigger.Once, which ignored the cap and read everything in one batch.

Parallelism and consumer pools

By default a batch has one task per Kafka partition with new data, so a 12-partition topic never uses more than 12 cores on the read side, however many executors you have. minPartitions asks the driver to split large offset ranges into several tasks. It is a hint, not an exact count, and it has a cost: two tasks reading the same partition at once need two consumers, so the consumer pool cannot reuse the cached one and fetch efficiency drops. Use it when partitions are few and processing per record is heavy, not as a substitute for adding Kafka partitions.

The executor-side pools are configured with spark.kafka.consumer.cache.capacity (default 64, a soft limit), spark.kafka.consumer.cache.timeout (5m idle) and matching settings under spark.kafka.consumer.fetchedData.cache.*. If one executor reads more than 64 distinct partitions, consumers are evicted and recreated every batch, which shows up as slow first polls. Raise the capacity in that case.

Worked example: sizing catch-up after an outage

A topic has 24 partitions and receives 30,000 records per second at about 1 KB each. The query triggers every 30 seconds, so a steady batch covers about 900,000 offsets, roughly 0.9 GB.

  • Measure: in a test, the job processes a 900,000-offset batch in 12 seconds. It has headroom of about 2.5x.
  • Cap: set maxOffsetsPerTrigger to 2,000,000. A catch-up batch then takes about 27 seconds, still inside the 30-second trigger, and memory per task stays bounded (about 83,000 records per partition).
  • Outage: the job is down for two hours, building a backlog of 216 million offsets. Each catch-up batch consumes 2 million while 0.9 million new offsets arrive, so the backlog shrinks by about 1.1 million per batch: 196 batches, roughly 1 hour 40 minutes to catch up.
  • Retention: if the topic keeps 24 hours of data, the job can survive an outage up to roughly 22 hours before failOnDataLoss fires because the checkpointed offsets have been deleted.

Watch query.lastProgress: the Kafka source reports start, end and latest offsets per partition, and metrics such as avgOffsetsBehindLatest and maxOffsetsBehindLatest. Export those to your metrics system; they are the real lag signal, since nothing is written to Kafka's consumer-group offsets.

p = query.lastProgress
src = p["sources"][0]
print(src["endOffset"])                       # {"orders": {"0": 1234, ...}}
print(src["metrics"].get("maxOffsetsBehindLatest"))
print(p["batchDuration"], p["numInputRows"])

Bounded reads for backfills and reconciliation

The same source also works with spark.read for bounded jobs, which is the right tool for a backfill or for checking what a stream should have produced. In batch mode startingOffsets defaults to earliest (latest is not allowed), and endingOffsets, endingOffsetsByTimestamp or endingTimestamp bound the read. A batch read never touches a streaming query's checkpoint, so it is safe to run alongside production.

day = (spark.read.format("kafka")
       .option("kafka.bootstrap.servers", "b1:9093,b2:9093")
       .option("subscribe", "orders")
       .option("startingTimestamp", "1790899200000")   # 2026-10-02 00:00 UTC, ms
       .option("endingTimestamp",   "1790985600000")   # 24 hours later
       .load())

expected = day.count()
written = spark.table("bronze.orders").where("ts >= '2026-10-02' AND ts < '2026-10-03'").count()
print(expected, written)

A reconciliation like this, comparing the records Kafka holds for a window with the records the stream wrote, catches silent loss from a mistaken failOnDataLoss=false or a sink bug. Note that timestamps here are Kafka record timestamps, which depend on the producer's message.timestamp.type setting, not the event time inside the payload, so allow for small differences at window edges.

When a streaming query must be restarted from a known point, for example after a bad deploy, the batch API is also how you find the offsets to start from: read the partitions around the incident, pick per-partition offsets, and pass them as the JSON startingOffsets of a new query with a new checkpoint.

Failure modes

  • Offsets out of range. Retention deleted data the checkpoint still points at. With failOnDataLoss=true the query stops, which is correct for financial data. Setting it to false makes the source skip ahead with a warning; do that only where loss is acceptable and you alert on it.
  • Deleted or recreated topic. A topic recreated with the same name restarts its offsets at 0, below the checkpoint. This is detected as possible data loss. Start a new checkpoint deliberately rather than turning the check off.
  • Aborted transactions read as data. Producers using transactions need consumers with kafka.isolation.level=read_committed; the Kafka client default reads uncommitted records.
  • Poll timeouts. TimeoutException after kafkaConsumer.pollTimeoutMs usually means a broker or network problem, or TLS handshakes failing. Fix the cause; a longer timeout just makes batches hang longer.
  • Unbounded first batch. A new query with startingOffsets=earliest and no cap tries to read the whole topic in one batch. Always set maxOffsetsPerTrigger on topics with long retention.
  • Changing the topic list. Some changes between restarts are allowed and some are not; check the checkpoint compatibility rules in the architecture article before editing a running query's source.

Trade-offs

The source's design is a trade: by owning offsets in the checkpoint instead of in Kafka, Spark gets deterministic replays and exactly-once with idempotent sinks, but loses Kafka-native lag monitoring and consumer-group tooling. Admission control trades latency for stability. minPartitions trades fetch efficiency for parallelism. failOnDataLoss trades availability for correctness. Make each choice explicitly, per pipeline, and write it down next to the query.

A useful rule of thumb: a pipeline feeding money, billing or compliance data keeps failOnDataLoss=true, a tight maxOffsetsPerTrigger and a nightly batch reconciliation; a pipeline feeding dashboards or feature freshness can accept skipping ahead after an outage, provided the skip raises an alert and someone decides whether to backfill. If several queries read the same topic, give each its own checkpoint and generated group id so that a restart or reset of one never moves another.

What to do next

  1. List every streaming query that reads Kafka and record its checkpoint path, startingOffsets, maxOffsetsPerTrigger and failOnDataLoss.
  2. Add a cap to any query without one, sized from a measured batch duration as in the worked example.
  3. Export maxOffsetsBehindLatest from lastProgress and alert when lag implies you will cross topic retention.
  4. Set kafka.isolation.level=read_committed wherever producers are transactional.
  5. Keep partition and offset in the first table you write; read foreachBatch in depth for idempotent writes, watermarks for event time, and streaming state for stateful operators downstream.
Key takeaway: The Kafka source plans each batch's offsets on the driver, records them in the checkpoint, and reads them with pooled consumers on executors. Starting options apply only to a fresh checkpoint, maxOffsetsPerTrigger keeps catch-up batches bounded, and because offsets are not committed to Kafka you must export lag from query progress yourself.