Case · Aviation · Anonymous client

From deeply nested XML events to nine analysis-ready business entities, with safe reprocessing, exact reconciliation and idempotency at every boundary.

Azure DatabricksLakeflowDelta LakeUnity CatalogPySpark

The context

Data arriving days after the decision.

An airline continuously received ticket lifecycle events through a GDS. Issuances, exchanges, refunds, cancellations, taxes and payment methods arrived through a message queue as deeply nested XML.

Consolidation relied on a handcrafted Spark Streaming implementation. Manual writes, fixed triggers and rules distributed across notebooks made the system hard to test and reprocess. Commercial analytics received consolidated information with days of latency.

Our mandate covered a ground-up replacement from the queue integration boundary to the business model. Before designing the architecture, we asked where low latency actually changed a decision. That question kept streaming complexity out of layers that did not need it.

Real-time only pays for itself when there is a decision on the other side that cannot wait for the nightly run.

The challenge

Replace the system without losing the original record.

The solution had to run continuously, handle complex XML, live with at-least-once delivery and preserve consistency in a model used for financial analysis. It also had to avoid turning every failure into manual intervention or creating an architecture that could not be promoted reliably across environments.

Our approach

Separate delivery, normalization and business rules.

We designed a serverless boundary to consume and preserve messages, a continuous pipeline for ingestion and normalization, and another pipeline to materialize the business model. Each layer received semantics suited to its problem, supported by versioned configuration and testable logic. The integration boundary became an observable contract between delivery and processing, while the original payload remained available for audit, diagnosis and replay.

Publishable architecture

An observable boundary between delivery and processing.

A message only leaves the queue after its original payload has been preserved. From there, two declarative pipelines turn the XML into normalized data and business entities.

Delivery01
  1. GDSXML events
  2. MQ queuelocked delivery
  3. Integrationidempotent ledger
acknowledged only after it is preserved
Preservation02
  1. Landingimmutable payload
the original is never rewritten
Modelling03
  1. Pipeline 1~40 tables
  2. Pipeline 29 entities
reprocessable at any time
Consumption04
  1. ConsumptionBI and analytics
Diagram redrawn for publication. It describes the shape of the solution without reproducing internal names, screens, identifiers or client resources.

Key decisions

Every engineering choice removed one kind of manual work.

01

Preserve before processing

Landing every message as a file looked like an extra hop. In practice, it created an immutable record, decoupled the pipeline from queue retention and made every row traceable to its source payload.

Result: manual reprocessing disappeared.

02

Idempotency at two boundaries

The ledger protects against redelivery of the same message after a failure. In the lakehouse, explicit business-key ordering deterministically selects the correct version of each fact.

Result: manual duplicate correction disappeared.

03

Fallback at document level

When a reference needed to allocate a value across segments was missing, the alternative method was applied to the whole document. The sum of the segments therefore remained equal to the document total.

Result: manual reconciliation disappeared.

04

Streaming where it changes a decision

Ingestion and normalization run continuously. The business model uses incremental materializations, avoiding fragile joins and aggregations across multiple streams for a marginal latency reduction.

Result: low latency with an operation people can understand and test.

In production

More than speed: a platform that can be operated with confidence.

~40normalized tables from one XML structure
9business entities materialized for analytics
4environments promoted from the same codebase
< 1 smedian integration component cycle in development measurements

What only appeared when the system met production

The technical article explains where we stopped streaming, how we handled redelivery and five failures that changed how we build pipelines.

Read the technical article