Spark Streaming, the API built on discretized streams (DStreams), was Spark's first streaming engine. It cuts an unbounded stream into small batches, turns each batch into an ordinary RDD, and runs a normal Spark job on it. That simple idea made streaming cheap to adopt in 2013 and still runs a surprising number of production pipelines in 2026.
It is also finished. The DStream API was formally deprecated in Spark 3.4.0 (SPARK-42075), and the current Spark Streaming Programming Guide opens by calling it a legacy project with no further updates, pointing users to Structured Streaming. You still need to understand it if you inherit a DStream job or must migrate one without losing its guarantees. This article explains the engine from first principles, walks a worked Kafka example, covers recovery, tuning and failure modes, and ends with a migration map and checklist.
The micro-batch model
Everything starts with a StreamingContext and a batch interval, say 10 seconds. Every 10 seconds the driver's job generator asks each input DStream for the data that arrived during that interval and wraps it as an RDD. Transformations such as map, filter and reduceByKey on a DStream are recorded in a graph, and at each tick the graph is replayed against the new RDD to produce a DAG of ordinary Spark stages. Nothing executes until an output operation such as foreachRDD or print is registered; a DStream graph with no output operation fails at start.
Two consequences follow from this design. First, latency is bounded below by the batch interval: a record that arrives just after a tick waits a full interval before it is even scheduled. Second, the engine is only stable if each batch finishes before the next one is due. If a batch takes 12 seconds on a 10-second interval, the next batch queues behind it, the scheduling delay grows by two seconds every batch, and the job falls further behind until it runs out of memory or someone notices. The Streaming tab in the Spark UI shows processing time and scheduling delay per batch; those two numbers are the health of a DStream job.
Receivers and direct streams
Data enters in one of two ways, and the difference matters for cores, parallelism and recovery.
Receivers. Socket, Kinesis and custom sources run a receiver: a long-running task pinned to one executor core. The receiver chops incoming records into blocks every spark.streaming.blockInterval (200 ms by default; the guide recommends not going below about 50 ms) and replicates them to another executor, since received input defaults to storage level MEMORY_AND_DISK_SER_2. Each block becomes one partition of the batch RDD, so a 10-second batch from one receiver has 10 / 0.2 = 50 partitions. Because every receiver permanently holds a core, the application needs more cores than receivers; in local mode use local[n] with n greater than the receiver count, or the job receives data and never processes it.
Direct streams. The Kafka 0-10 integration does not use a receiver. At each tick the driver asks Kafka for the latest offsets, decides an offset range per topic partition, and creates an RDD whose partitions map one-to-one onto Kafka partitions. Executors read their ranges directly. There is no receiver core, no block replication and no write-ahead log; the offset ranges are the source of truth, and replaying a batch means re-reading the same ranges. This is why almost every production DStream job on Kafka uses the direct stream.
A receiver that has acknowledged data to its source can lose it if the executor dies before the block is processed. Setting spark.streaming.receiver.writeAheadLog.enable to true makes receivers write blocks to a write-ahead log in the checkpoint directory first, which gives at-least-once delivery at the cost of an extra durable write per block.
Worked example: a Kafka click counter
Here is a complete job in the shape most inherited pipelines have: read click events from Kafka, keep a running count per page with timeouts, compute a five-minute sliding total, and write results to an external store with offsets committed only after the write. It is Scala because the Kafka direct stream and mapWithState are Scala and Java APIs.
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.SparkConf
import org.apache.spark.streaming._
import org.apache.spark.streaming.kafka010._
import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
object ClickCounts {
val checkpointDir = "hdfs:///checkpoints/click-counts-v3"
def createContext(): StreamingContext = {
val conf = new SparkConf().setAppName("click-counts")
.set("spark.streaming.backpressure.enabled", "true")
.set("spark.streaming.kafka.maxRatePerPartition", "2000")
.set("spark.streaming.stopGracefullyOnShutdown", "true")
val ssc = new StreamingContext(conf, Seconds(10))
ssc.checkpoint(checkpointDir)
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "broker1:9092,broker2:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "click-counts",
"auto.offset.reset" -> "latest",
"enable.auto.commit" -> (false: java.lang.Boolean))
val stream = KafkaUtils.createDirectStream[String, String](
ssc, PreferConsistent, Subscribe[String, String](Seq("clicks"), kafkaParams))
// Offsets are only visible on the RDD that came straight from the stream.
stream.foreachRDD { rdd =>
val ranges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
val pages = rdd.map(r => (pageOf(r.value), 1L))
pages.reduceByKey(_ + _).foreachPartition { part =>
val sink = SinkPool.borrow() // one connection per partition, not per record
try part.foreach { case (page, n) => sink.addCount(page, n) }
finally SinkPool.giveBack(sink)
}
stream.asInstanceOf[CanCommitOffsets].commitAsync(ranges)
}
// Running total per page; idle keys expire after 30 minutes.
val spec = StateSpec.function((page: String, n: Option[Long], state: State[Long]) => {
if (state.isTimingOut()) (page, state.get())
else { val total = state.getOption().getOrElse(0L) + n.getOrElse(0L); state.update(total); (page, total) }
}).timeout(Minutes(30))
stream.map(r => (pageOf(r.value), 1L)).mapWithState(spec).print()
// Five-minute window sliding every 10 s, using an inverse function.
stream.map(r => (pageOf(r.value), 1L))
.reduceByKeyAndWindow(_ + _, _ - _, Minutes(5), Seconds(10))
.checkpoint(Seconds(60))
.print()
ssc
}
def main(args: Array[String]): Unit = {
val ssc = StreamingContext.getOrCreate(checkpointDir, createContext _)
ssc.start()
ssc.awaitTermination()
}
}Three details in that code are the ones teams get wrong. The cast to HasOffsetRanges works only on the RDD produced directly by createDirectStream; after a map or shuffle the partitions no longer correspond to Kafka partitions and the cast fails. The commit happens after the output, so a crash between output and commit replays the batch: delivery is at-least-once, and the sink must tolerate duplicates. And the whole graph is built inside createContext, because getOrCreate only calls that function when no checkpoint exists; on restart the graph is deserialized from the checkpoint instead.
Windows and state
DStreams offer three kinds of state. Windowed operations such as window, countByWindow and reduceByKeyAndWindow combine the RDDs of the last N batches; window length and slide must both be multiples of the batch interval. The inverse-function form shown above is incremental: each slide adds the newest batch and subtracts the batch that left the window, instead of re-reducing five minutes of data every ten seconds. It requires checkpointing, because otherwise the lineage of the running window grows forever.
updateStateByKey is the original stateful operator. It calls your update function for every key in the state on every batch, whether or not the key received new data, so its cost grows with total state size rather than with input rate. mapWithState (Scala and Java only) touches only keys with new data or keys that are timing out, supports per-key timeouts, and is the better choice for any large keyspace. Both keep state as RDDs that must be checkpointed periodically.
The guide's advice for the state checkpoint interval is five to ten sliding intervals, with a default that is a multiple of the batch interval and at least 10 seconds. Checkpointing every batch makes HDFS or S3 writes dominate processing time; checkpointing too rarely makes lineage, and therefore recovery time and task size, grow.
Checkpoints, recovery and the upgrade trap
The checkpoint directory holds two different things. Metadata checkpoints store the configuration, the DStream graph, and the list of batches that were queued but not completed; they let a restarted driver pick up where it left off. Data checkpoints store generated state RDDs so that stateful operators do not need to recompute from the beginning of time.
Driver recovery works like this: the driver dies, a cluster manager with supervision (or your own wrapper) restarts it, StreamingContext.getOrCreate finds the checkpoint, deserializes the graph, and re-runs the incomplete batches. For a direct Kafka stream the incomplete batches carry their offset ranges, so they read exactly the same data again.
The trap is the word deserializes. The checkpoint contains serialized Java objects from your application classes. The guide states plainly that after you change and redeploy the application code, the old checkpoint cannot be restored; the usual result is a deserialization error at startup, or, worse, a job that starts with subtly stale configuration. Every code upgrade therefore needs a plan: run the new version with a new checkpoint directory (the -v3 suffix in the example), start it from offsets you control, and either let the old job drain with a graceful stop or run both in parallel and switch output. Because Kafka offsets were committed back to the consumer group, the new job can resume from the committed position, but any mapWithState or window state is lost and must be rebuilt or reloaded.
Rate limits, backpressure and tuning
Tuning a DStream job is mostly about keeping processing time under the batch interval with headroom. Start from the arithmetic. A topic with 24 partitions, a 10-second batch and spark.streaming.kafka.maxRatePerPartition = 2000 caps each batch at 24 x 2,000 x 10 = 480,000 records. If the job processes 60,000 records per second when healthy, a 480,000-record batch takes about 8 seconds, which leaves 20 percent headroom. After an outage the backlog is drained at the cap, never in one giant batch that would take minutes and stall everything behind it.
Backpressure (spark.streaming.backpressure.enabled) adds a feedback controller that estimates the sustainable rate from recent batch timings and lowers the ingestion rate when batches slow down. Keep the static cap as well: backpressure reacts to past batches, so the first batch after a restart is bounded only by the static limit. For receivers the equivalent cap is spark.streaming.receiver.maxRate.
| Symptom | Likely cause | First fix |
|---|---|---|
| Scheduling delay climbs steadily | Processing time > batch interval | Raise parallelism or rate cap; lengthen interval |
| Few long tasks per batch | Receiver blocks too large or few Kafka partitions | Lower blockInterval, add partitions or repartition |
| Batch time spikes every N batches | State or window checkpoint | Increase checkpoint interval to 5-10 slides |
| Executors run out of memory | updateStateByKey over a huge keyspace | Switch to mapWithState with timeouts |
| Sink saturated with connections | Connection opened per record | Open per partition from a pool |
Failure modes
- Silent backlog. Nobody alerts on scheduling delay, so a job that slowed down weeks ago is now hours behind. Alert when scheduling delay exceeds one or two batch intervals.
- Checkpoint incompatibility after deploy. The new build cannot read the old checkpoint. Version the directory and script the cut-over.
- Duplicates in the sink. A failure between output and offset commit replays the batch. Make writes idempotent (upserts keyed by a natural id) or store offsets in the same transaction as the results.
- Lost receiver data. Without the receiver WAL, an executor crash loses buffered blocks that the source already considers delivered.
- Graceful stop that is not graceful. Killing the driver mid-batch forces a replay. Use
spark.streaming.stopGracefullyOnShutdownor callstop(stopSparkContext = true, stopGracefully = true)so queued batches finish first. - Wall-clock windows on event-time questions. DStream windows are defined by arrival batches, not event timestamps, so late data lands in the wrong window and there is no watermark to bound it.
Migrating to Structured Streaming
The last failure mode is the main reason to migrate. Structured Streaming treats the stream as an unbounded table, supports event-time windows with watermarks, has end-to-end exactly-once for supported sinks, and keeps offsets in its own checkpoint in a format designed to survive many code changes. Most DStream constructs have a direct counterpart:
| DStream | Structured Streaming |
|---|---|
| StreamingContext + batch interval | SparkSession + trigger (processingTime) |
| KafkaUtils.createDirectStream | spark.readStream.format("kafka") |
| reduceByKeyAndWindow (arrival time) | groupBy(window(eventTime, ...)) + withWatermark |
| mapWithState / updateStateByKey | flatMapGroupsWithState or transformWithState |
| foreachRDD + foreachPartition | foreachBatch or a custom ForeachWriter |
| commitAsync of offset ranges | offsets tracked in the query checkpoint |
| Checkpoint tied to compiled classes | Checkpoint survives many query changes |
Run the migration as a side-by-side comparison. Start the new query on a separate consumer group reading the same topic, write to a shadow table, and diff aggregates per window for a few days. Differences usually come from the semantic change from arrival time to event time, which is the point of migrating, so decide which number is correct before cutting over. The Structured Streaming overview, Kafka with Structured Streaming, watermarks, stateful processing and foreachBatch pages cover the destination side in detail.
Trade-offs
Keeping a DStream job is defensible when it is stable, its semantics are arrival-time by design, nobody needs to change its code, and its Spark version is still within your support window. The cost is that every change is risky (checkpoint incompatibility), there will be no fixes for new problems, and the integrations that remain are frozen at their current state. Migrating costs engineering time and a careful parity check, and it buys event-time correctness, watermark-bounded state, a supported engine and a lower-risk upgrade path. For anything that is still evolving, migrate; for a frozen job that is scheduled to retire, pin the Spark version and monitor it closely.
What to do next
- Inventory every job that imports
org.apache.spark.streamingand record its source, batch interval, stateful operators and sink. - Add alerts on scheduling delay and processing time per batch from the Spark UI metrics.
- Confirm each Kafka job uses the direct stream, commits offsets after output, and writes idempotently.
- Set a static rate cap and turn on backpressure, then test a restart against a large backlog.
- Version the checkpoint directory and write down the cut-over procedure for the next code change.
- Pick one job, rebuild it in Structured Streaming, and run both side by side against a shadow sink until the outputs agree or the differences are explained.