Where to stop streaming, how to build idempotency at every boundary and five lessons from an airline ticketing platform in production.

Data engineering · Azure Databricks · Technical read

In this article
  1. 01 The decision comes before the architecture
  2. 02 The extra hop that bought reprocessability
  3. 03 Idempotency is not a box checked once
  4. 04 Where we stopped streaming
  5. 05 Why Bronze became a table
  6. 06 Incremental ingestion without another event service
  7. 07 Aggregate consistency over individual precision
  8. 08 The INNER JOIN that returned zero rows
  9. 09 Deduplication that changed between runs
  10. 10 The official standard that duplicated a dimension
  11. 11 The silent debt of small files
  12. 12 Retiring legacy is engineering too
  13. 13 The framework stays at the edge of the logic
  14. 14 What turns a pipeline into a platform
  15. 15 How to discuss latency without inventing a number
  16. 16 The outcome that matters
The requirement was real-time. The first thing we did was not architecture. We asked why.
01

The decision comes before the architecture

Continuous processing has a permanent cost: active compute, operational complexity and a class of failures that only exists in streaming. The architecture only pays for itself when a real decision loses value by waiting until the next day.

Commercial reporting alone did not justify streaming throughout the solution. The ability to detect anomalous issuance patterns while they are still actionable changes that equation. We used this criterion to decide where low latency was necessary and, just as deliberately, where it was not.

That question became a design constraint. We measured value through the latency of the complete system, not through the number of boxes labeled streaming in a diagram.

A real-time architecture must begin with the decision that cannot wait.

02

The extra hop that bought reprocessability

The obvious alternative was to consume the queue directly from Spark Structured Streaming. We rejected it and materialized every message as a file in the landing zone before analytical processing began.

A queue has limited retention. The file preserves exactly what the source sent and lets us fix a parser months later without asking for another extract. It also decouples pipeline availability from queue behavior and creates a simple operational boundary: did the file arrive or not?

In a model used for financial analysis, lineage and auditability are not conveniences. Every row must trace back to the payload that produced it. Saving one network hop was worth less than guaranteed replay and diagnosis.

In financial pipelines, reprocessability is worth more than shaving seconds off latency.

03

Idempotency is not a box checked once

End-to-end delivery is not magically exactly-once. The design combines at-least-once delivery with consumer-side idempotency at two independent boundaries.

At the integration boundary, a message remains locked until the write finishes. If a batch fails halfway through, it returns to the queue. A ledger records only the items already written on that failure path and prevents duplicate files. After a successful batch, that temporary state is removed.

Inside the lakehouse, the source may send legitimate revisions of the same business fact. Deduplication therefore uses a business key and explicit ordering to select the right version. Technical redelivery and a business revision look similar but require different defenses.

Idempotency must exist at every boundary where delivery can repeat.

04

Where we stopped streaming

Ingestion and normalization run continuously. The business model is materialized incrementally. This boundary does not make the system less ambitious. It protects correctness.

The final entities require aggregations and joins across several sources. Forcing streaming semantics into that layer would introduce watermarks across multiple flows, incompatible output modes and behavior that is difficult to explain during a failure. The marginal latency gain was not worth it.

Near real-time is an end-to-end property. The right boundary delivers the business window without importing complexity into layers that do not benefit from it.

Choosing where to stop streaming is an engineering decision, not a lack of ambition.

05

Why Bronze became a table

The Bronze layer fed dozens of downstream tables. Keeping it as a view looked lighter, but every consumer would open an independent reader against the same source, with separate checkpoints and no reliable read order.

Materializing it created one ingestion checkpoint and consistent fan-out. The cost of persisting that boundary bought predictability for every transformation that followed.

A decision that looks like performance can actually be a correctness decision.

06

Incremental ingestion without another event service

We used file notifications managed by the governed storage layer. That removed an additional event service, its queue, credentials and the risk of configuration drift across four environments.

The schema was declared explicitly because inference over nested XML in streaming is unstable. Out-of-schema records are preserved in a rescue column, while lineage metadata enters at ingestion time. Bronze never discards a payload just because it cannot interpret it yet.

stream = (
    spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "xml")
        .option("cloudFiles.useManagedFileEvents", "true")
        .option("rescuedDataColumn", "_rescued")
        .schema(explicit_schema)
        .load(landing_zone_path)
)

Less infrastructure also means fewer credentials, less drift and fewer failure points.

07

Aggregate consistency over individual precision

One entity allocated a document value across its segments proportionally. The reference used in the calculation did not cover every possible pair. Applying a fallback only to the segment with a missing reference made the final sum diverge from the document total.

We changed the unit of the rule. If one reference was missing, the entire document switched to the alternative method. Individual segment precision could decrease, but document-level reconciliation remained exact in every case.

The change removed manual reconciliation because the rule now operated at the same level where the business verifies consistency.

Fallback must be designed at the reconciliation unit, not at the record with missing data.

08

The INNER JOIN that returned zero rows

One entity joined documents to the source considered primary for value and currency. In production, the result contained zero rows. For that event type, the information never appeared in that source. It existed only in an alternative.

The fix was a LEFT JOIN with an explicit coalesce hierarchy across sources. The issue was not syntax. The join claimed that every document had a mandatory match, a business assumption hidden in ordinary SQL.

After the fix, we began reviewing joins as contracts for cardinality and presence. Before choosing the type, we ask what absence means in the domain and whether dropping a row is truly valid.

INNER or LEFT is not a style preference. It is a testable claim about the domain.

09

Deduplication that changed between runs

The first version used dropDuplicates. The operation removes duplicates but does not guarantee which row survives when two versions share a key. Reprocessing the same window could produce a different result even when no input had changed.

We replaced it with a window partitioned by business key and explicitly ordered by event time. The selected version no longer depended on the physical distribution of the DataFrame.

This detail is central to replay. A reprocessable platform cannot merely finish without an error. It must return the same answer for the same input, including after a cluster or partitioning change.

window = Window.partitionBy("business_key").orderBy(
    F.col("event_timestamp").desc()
)

deduped = (
    events.withColumn("_row_number", F.row_number().over(window))
          .filter(F.col("_row_number") == 1)
          .drop("_row_number")
)

Idempotency is the difference between replaying confidently and being afraid to run the pipeline again.

10

The official standard that duplicated a dimension

A dimension was enriched with geographic hierarchy from a reference based on an international standard. The join appeared naturally one-to-one, yet it increased the row count.

The reference preserved the history of codes assigned to different entities over time. Official did not mean unique. We filtered historical entries and validated key cardinality before the join.

We also adopted a convention: a DataFrame only earns the lookup suffix after it has been reduced and validated for the expected cardinality. The name communicates a guarantee, not merely an intention.

Standards carry history. Validate cardinality before assuming a LEFT JOIN preserves rows.

11

The silent debt of small files

In batch, poor partitioning leaves one bad set of files. In continuous execution, every micro-batch repeats the problem. A table starts healthy and slowly degrades until listing and opening files dominates reads.

We tuned adaptive execution, partition coalescing, target size and automatic optimization. The goal was not a perfect number for every batch, but preventing the debt from growing by itself for weeks.

Monitoring job duration alone does not catch this pattern early. File count and size distribution are part of the operational health of a continuously updated table.

In streaming, small files are not a one-off nuisance. They are debt that compounds by itself.

12

Retiring legacy is engineering too

The migration did not end when the new pipeline produced the right result. Dozens of manually managed external objects had to leave the catalog before the new operation, without deleting underlying data or touching tables already promoted.

We created a versioned cleanup process with explicit targets and a protection list. Retirement became reviewable and repeatable like any other artifact instead of relying on a manual command session.

The process ended coexistence without data loss or an availability window. The new operation only became simple after the old one stopped competing for names, responsibilities and attention.

A migration ends when legacy leaves the stage in an auditable way.

13

The framework stays at the edge of the logic

Declarative pipelines do not need to bind business logic to the framework runtime. We kept decorators in thin wrappers that only resolve sources. Transformations live in pure functions that receive and return DataFrames.

This pattern lets us test rules with synthetic data and pytest without starting the full pipeline. We do not claim complete test coverage. The demonstrable value is structural: each rule can be exercised in isolation, while the same code moves through development, quality, pre-production and production through configuration alone.

@pipeline.table(name="business_entity")
def business_entity():
    return build_business_entity(
        spark.table(source_a),
        spark.table(source_b),
    )

def build_business_entity(df_a, df_b):
    # All business logic lives in this pure function.
    ...

Frameworks should orchestrate logic, not hide it.

14

What turns a pipeline into a platform

Pipelines are defined in versioned configuration, with catalogs, paths and compute policies resolved by environment. The same transformation file crosses four environments without local names or conditional branches. That narrows the distance between what was tested and what reaches production.

Observability also uses native platform artifacts. The pipeline event log is materialized in the catalog and queried through SQL. Execution time, counts, updates and lineage remain available without scattering logging calls through every transformation.

This choice does not remove alerts or operational responsibility. It creates one investigation source. When a flow is late, the on-call engineer can separate a delivery failure, ingestion delay and modeling issue without opening several systems before forming a hypothesis.

Naming conventions, reproducible deployment and a clear responsibility boundary rarely appear in a demo. Yet these decisions determine whether another team can maintain the platform after the project ends.

Operability begins when the on-call engineer can understand what happened without depending on the person who wrote the pipeline.

15

How to discuss latency without inventing a number

Complete latency has three parts: message availability to file, file to Bronze and Bronze to a materialized entity. We measured only the first in a transferable way. Across two independent development workflows, ledger lookup, writing and confirmation completed in under one second at the median.

That number demonstrates that the stateless component has a stable execution cost. It does not demonstrate end-to-end latency, production throughput or message frequency. Adding estimates for the other legs would produce a commercially attractive but technically indefensible figure.

The correct way to complete the measurement is to compare file modification time with the Bronze ingestion timestamp, then use the event log to measure entity materialization. A published result should include p50, p95, sample size, period and environment.

Until that sample exists, we describe the system as near real-time and keep the sub-second metric scoped to the component actually observed. Precise language is part of engineering, especially when technical buyers will read the page.

A metric belongs in a case only when its source, scope and limitations can be explained.

16

The outcome that matters

The system runs in production and continuously provides information that previously reached analytics with days of latency. The same codebase crosses four environments without changes to the business logic.

The immutable landing zone removed handcrafted reprocessing. Idempotency at two boundaries removed manual duplicate correction. Document-level fallback removed manual reconciliation. These are results without invented ROI and with a direct relationship between an engineering decision and the operational work it removed.

The only safe performance measurement covers the integration component: ledger lookup, write and queue confirmation complete in under one second at the median in development histories. We do not turn that number into end-to-end latency. The complete system is honestly described as near real-time.

A platform survives when replay, auditability and operation are design properties, not emergency procedures.

The complete architecture, without marketing shortcuts

The case brings together the context, sanitized design, confirmed outcomes and the connection between each engineering decision and the manual work it removed.

View the complete case