The event arrived two days late. How to handle late data in Databricks

Design for lateness instead of arguing with it

The order was placed on Monday. It reached the warehouse system on Monday. It reached our tables on Wednesday.

Nobody did anything wrong. A mobile app was offline in a train tunnel. A partner file was delivered a day late. A queue backed up for two hours during a deploy. The event was real the whole time. It just showed up after we had already counted Monday and told everyone the number.

This is late-arriving data, and it is one of the quietest sources of distrust in a data platform. The pipeline runs green. The dashboard changes anyway.

Two clocks, not one

Every event carries two timestamps, and confusing them is where most late data problems start.

Event time is when the thing happened. The customer clicked at 21:14.

Processing time is when your pipeline saw it. That could be 21:15, or Wednesday.

If you group by processing time, your numbers are easy to compute and impossible to explain. Monday's revenue will keep changing shape depending on when files arrived. If you group by event time, your numbers mean what people think they mean, and you accept a new job: the past can still receive data.

Late data is not a bug to eliminate. It is a normal property of distributed systems that your design has to expect.

Open a raw table and look for yourself before you build anything on top of it. This single query tells you how late your data actually runs.

SELECT
  DATE_DIFF(DATE(ingested_at), DATE(event_time)) AS days_late,
  COUNT(*) AS event_count
FROM workspace.default.bronze_events
GROUP BY days_late
ORDER BY days_late

Most teams are surprised here. They expected everything to land the same day, and they find a small tail stretching out three or four days. That tail is your real requirement.

Step one: keep both timestamps in bronze

The bronze layer should record what arrived and when it arrived, without judgement. Add the ingestion timestamp as you read the files, and never overwrite the original event time.

from pyspark.sql import functions as F

raw_events = (
    spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.schemaLocation", "/Volumes/workspace/default/book_data/_schema/events")
        .load("/Volumes/workspace/default/book_data/events/")
)

bronze_events = raw_events.withColumn("ingested_at", F.current_timestamp())

(
    bronze_events.writeStream
        .option("checkpointLocation", "/Volumes/workspace/default/book_data/_checkpoints/bronze_events")
        .trigger(availableNow=True)
        .toTable("workspace.default.bronze_events")
)

Two columns, two clocks. Now every question about lateness has an answer in the data instead of in someone's memory.

Step two: choose a lateness budget out loud

A lateness budget is a decision, not a setting. It answers one question: how long do we keep a day open before we stop accepting new events for it?

Pick the number with the people who read the reports. Three days is a common answer. Once you have it, write it down somewhere visible, because it changes what everyone can promise. A two-hour budget means yesterday is final by breakfast. A seven-day budget means last week can still move.

In streaming, this budget becomes a watermark. The watermark tells Spark how far behind the newest event time it should keep waiting.

windowed_counts = (
    spark.readStream.table("workspace.default.bronze_events")
        .withWatermark("event_time", "3 days")
        .groupBy(F.window("event_time", "1 day"), "country")
        .agg(F.count("*").alias("event_count"))
)

Watermarks are how a streaming job stays cheap. Without one, Spark has to hold state for every window forever, in case something ancient arrives. With one, it can close old windows and release memory. The trade is explicit: events later than the budget are dropped from that aggregate.

Dropped silently is the part that bites people, so measure it.

SELECT COUNT(*) AS beyond_budget
FROM workspace.default.bronze_events
WHERE ingested_at > event_time + INTERVAL 3 DAYS

If that count is regularly above zero, your budget is wrong, not your data.

Step three: make the silver layer restate, not append

Batch pipelines usually get late data wrong in the same way. They append yesterday's file to yesterday's totals, and the totals drift.

The fix is to recompute the affected days rather than add to them. Find which event days appeared in this run, then rewrite exactly those days.

new_rows = spark.read.table("workspace.default.bronze_events").filter(
    F.col("ingested_at") >= F.lit("2026-08-19")
)

affected_days = [row.event_date for row in
                 new_rows.select(F.to_date("event_time").alias("event_date")).distinct().collect()]

day_list = ", ".join([f"'{day.isoformat()}'" for day in affected_days])

daily_totals = (
    spark.read.table("workspace.default.bronze_events")
        .withColumn("event_date", F.to_date("event_time"))
        .filter(F.col("event_date").isin(affected_days))
        .groupBy("event_date", "country")
        .agg(F.count("*").alias("event_count"))
)

(
    daily_totals.write
        .mode("overwrite")
        .option("replaceWhere", f"event_date IN ({day_list})")
        .saveAsTable("workspace.default.silver_daily_events")
)

replaceWhere is the important part. It replaces one clean slice of the table in a single atomic commit, so readers never see a half-updated day. This is the same discipline we use in safe reruns and in a full backfill. Late data is really just a small backfill that happens every day.

Step four: tell people which days are still open

The technical fix is only half the work. The other half is expectation.

Publish a small table that says, for each event day, whether it is still open, and when it was last recomputed. A dashboard can read it and show a quiet note: Monday is still open and may change until Thursday.

CREATE OR REPLACE TABLE workspace.default.event_day_status AS
SELECT
  DATE(event_time) AS event_date,
  MAX(ingested_at) AS last_arrival,
  CASE
    WHEN DATE(event_time) >= CURRENT_DATE() - INTERVAL 3 DAYS THEN 'open'
    ELSE 'final'
  END AS status
FROM workspace.default.bronze_events
GROUP BY DATE(event_time)

One table, and a hundred confused messages disappear. People do not mind numbers that move. They mind numbers that move without warning.

Step five: keep an undo

Even a careful restatement can go wrong. Delta Lake keeps history, so you always have a way back.

DESCRIBE HISTORY workspace.default.silver_daily_events;

RESTORE TABLE workspace.default.silver_daily_events TO VERSION AS OF 41;

Knowing you can undo is what makes it possible to fix things calmly during the day instead of nervously at night.

What to remember

Late data is not a failure. It is the honest behaviour of systems made of phones, networks, partners, and deploys.

Design for it with four decisions. Keep event time and ingestion time. Agree a lateness budget with the people who read the numbers. Restate whole days instead of appending to them. Say clearly which days are still open.

Do that, and the same event arriving two days late becomes a routine correction rather than a crisis.

Continue learning