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

Lambda Architecture combines three responsibilities: a batch layer recomputes results from retained history, a speed layer processes new events with low latency, and a serving layer makes the results available for queries. Apache Spark can run both processing paths: scheduled Spark SQL or DataFrame jobs for historical computation, and Spark Structured Streaming for incremental updates.

How the three layers work

Layer Responsibility Typical Spark role What it publishes
Batch Recompute results from the complete retained event history, including corrections to earlier data. Scheduled Spark SQL or DataFrame jobs. Authoritative historical views or tables.
Speed Process recent events incrementally so fresh results are available before the next batch recomputation. Spark Structured Streaming transformations. Recent or in-progress results, often with state for windows, joins, or deduplication.
Serving Expose queryable views that combine or reconcile batch and speed outputs. Spark may prepare outputs; the serving technology depends on query needs. Tables, operational databases, search indexes, dashboards, or APIs.

AWS describes Lambda as mixing batch and stream processing and making the combined data available through a serving layer. In practice, events commonly enter through Kafka, Amazon Kinesis, or another message bus. Retaining an immutable or append-oriented source history gives batch jobs durable input for replay and recomputation.

How to implement the pattern with Spark

  1. Ingest and retain events. Read events from a message bus such as Kafka or Kinesis, and preserve the source history in durable storage. Treat the retained history as the input for replay and full recomputation, rather than relying only on the stream processor’s current state.
  2. Build the batch path. Schedule Spark SQL or DataFrame jobs to read the complete history, apply the agreed business transformations, correct prior errors where needed, and publish authoritative tables or views.
  3. Build the speed path. Use Spark Structured Streaming to read new events, apply incremental transformations, manage any required state, and write fresh results. Spark’s structured APIs use the DataFrame and Dataset programming model across batch and streaming, which can reduce the need to maintain separate technology stacks.
  4. Make the serving view explicit. Decide how query-facing results combine the authoritative batch output and the newer speed output. The right serving store and reconciliation method depend on query shape, freshness requirements, consistency needs, and scale.
  5. Plan for correction and replay. Define how a corrected historical result replaces or reconciles with speed-layer output, and how the stream or batch job can be restarted after a failure. Test that the query-facing result remains coherent during and after recomputation.

Databricks’ reference architecture describes Structured Streaming reading event queues such as Apache Kafka or AWS Kinesis, with downstream processing and serving systems. Its production guidance also identifies Pulsar, Pub/Sub, Delta change feeds, and Iceberg change feeds as possible low-latency sources. These are options, not requirements: choose sources and sinks that meet the workload’s replay, latency, and query needs.

What determines correctness and operating cost

Checkpoints, state, and recovery

Structured Streaming’s programming guide describes checkpointing and write-ahead logs supporting end-to-end exactly-once fault tolerance in its documented micro-batch model. Stateful aggregations, stream-to-stream joins, and deduplication depend on durable checkpoints and deliberate state management. Sink behavior also matters: make writes idempotent where necessary so retries do not create unintended duplicate business results.

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

Event time and late data

For event-time windows, joins, or deduplication, choose watermarks to define how long the stream retains state while waiting for late events. A watermark is a policy trade-off: waiting longer can accommodate later arrivals but can retain more state; discarding data beyond the chosen lateness boundary can affect results. Set the policy to match the data’s lateness pattern and the business cost of excluding late events.

Output mode and trigger behavior

Append, update, and complete output modes affect which results are written as processing advances. Trigger interval, input rate, state-store size, sink throughput, and available cluster capacity also influence freshness and cost. Validate these together under expected load; a short trigger interval alone does not guarantee low end-to-end latency.

Spark’s Structured Streaming documentation gives 100 milliseconds as an example of latency achievable with its default micro-batch engine. That is a documented lower-bound example, not a general service guarantee: actual latency depends on the trigger, workload, state, source and sink behavior, cluster capacity, and backpressure. Databricks separately documents real-time processing modes, so any latency target should identify the processing mode and workload rather than assume all Spark streaming runs behave alike.

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

Lambda or Kappa: which pattern fits?

Kappa removes Lambda’s separate batch-processing path and treats the stream as the primary computation. That can reduce duplicated logic when streams are replayable and the stream-processing approach meets the workload’s correctness and recovery needs. Lambda retains an independent historical recomputation path alongside low-latency updates, but keeping the two paths semantically aligned adds engineering and operational work.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Decision factor Lambda tends to fit when… Kappa tends to fit when…
Historical correction Recomputing authoritative results from retained history is important. Replaying the retained stream can produce the required corrected result.
Freshness Recent updates must be available before the next batch recomputation. The stream path alone can meet freshness and query requirements.
Logic and operations The separate batch path justifies its additional logic and operational burden. A single stream-oriented computation can reduce duplication without unacceptable replay or operational cost.
Data and recovery Retention, replay cost, or correctness requirements favor a distinct full-history process. Replayable streams and the stream processor’s guarantees are sufficient for recovery and recomputation.
Serving needs Combining authoritative history with newer results meets the required query behavior. The stream-produced result can be served directly in the needed form.

Choose by evaluating freshness and tail latency, replay and correction behavior, handling of late or out-of-order events, state size, serving-query requirements, infrastructure cost, and the operational cost of maintaining one versus two processing paths. Neither pattern is automatically more correct or less expensive; the answer depends on retention, replay cost, correctness needs, and the complexity the team can operate.

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.