ESC
Type to search guides, tutorials, and reference documentation.

Big Data Architecture

How large-scale data platforms are layered — ingestion, storage, table formats, processing and serving — with the tradeoffs, ownership models and failure modes of each layer.

A big data architecture is the set of decisions that determine how data enters a system, where it is stored, what physical shape it is stored in, what reads it, and who is allowed to change any of that. It is not a product, and it is rarely the diagram of vendor logos it gets drawn as. The durable part is the layering — ingestion, storage, table format, processing, serving — because every platform answers the same questions at each layer whether it runs on one server or a thousand. Most production pain traces back to one layer's decision leaking into another: a storage layout chosen for the writer that the readers cannot scan, or a transform welded to a source that can never be replayed.

The five layers

Ingestion

Ingestion is the boundary where data stops being someone else's problem and becomes yours. Two properties matter more than throughput: whether the source can be re-read, and whether a record can be identified. A re-readable source — an append-only log with retention, an object-store prefix, a table with a monotonic cursor — lets you recover from a bad deploy by replaying it. A source you cannot re-read, such as a webhook you already acknowledged or a queue you consumed destructively, forces the landing zone itself to be the durable record.

The resulting rule is boring and almost always right: land raw, unmodified and immutable, with the ingest timestamp attached, then transform downstream. Raw landing costs storage. An unreproducible transformation costs an incident, and then a migration.

Storage and file layout

Analytical workloads read few columns across many rows, so analytical storage is columnar: values of the same column sit together, compress well because they are homogeneous, and let an engine skip columns it was not asked for. Transactional storage is row-oriented for the opposite reason — it fetches whole records by key. Choosing the wrong orientation is not a tuning problem, it is an architecture problem, and it shows up as scans that read everything.

Within columnar storage, the two levers are partitioning and file size. Partitioning writes data into directories keyed by a column so that a filter on that column eliminates whole directories without opening them. It only helps if queries actually filter on the partition key, and it hurts when the key is too granular: partitioning by the second produces a directory tree the planner spends longer listing than the query spends scanning. File size is the mirror image — many tiny files mean per-file overhead (open, seek, read footer, schedule a task) dominates useful work. Compaction into larger files is maintenance every table-based platform eventually needs.

Table format and metadata

A directory of files is not a table. Without a metadata layer, two writers can half-commit, a reader can observe a partially written batch, and a schema change is whatever the last writer happened to emit. A table format adds an atomic commit over a file listing, snapshot isolation for readers, explicit schema evolution, and the ability to query the table as of a previous snapshot. Open implementations of this idea include Apache Iceberg, Delta Lake and Apache Hudi; they differ in their metadata structure and their update mechanics, but they exist to solve the same problem.

The strategic point is that engines are replaceable and formats are not. Query engines get swapped every few years; the petabytes in your object store do not get rewritten casually. Choose the storage and table format with more care than the compute.

Processing

Processing is where cost is actually incurred. The dominant cost in distributed processing is not CPU, it is data movement — the shuffle that redistributes rows across workers so that rows sharing a key land together. Anything that reduces shuffle (filtering early, pruning columns, pre-partitioning by join key, broadcasting a small side of a join) is worth more than anything that tunes the workers. The most common pathology is skew: one key holds a disproportionate share of rows, one task runs long after all others finish, and the cluster idles while a single worker grinds.

Serving

Serving is the read side, and the mistake is to have one of it. Point lookups by key, interactive aggregations over recent data, and large historical scans have contradictory optimal layouts. Trying to satisfy all three from a single table produces a table that is mediocre at each. A serving layer usually means a warehouse or query engine for analysis, a separate low-latency store for anything user-facing, and an explicit job that keeps the second derived from the first.

The reference shapes

Lambda runs a batch path and a speed path in parallel and merges them at read time. It is honest about the fact that batch reprocessing is the correctness backstop, but it makes you implement the same business logic twice in two runtimes, where the two implementations drift.

Kappa keeps only the streaming path and treats reprocessing as replaying the log through the same code. It removes the duplicate implementation but makes retention, replay throughput and state migration into first-class problems. See real-time analytics and stream processing for what that commits you to.

Layered (bronze/silver/gold, or raw/cleansed/curated) is less an architecture than a discipline: raw ingest is never edited, cleansing is idempotent and reproducible from raw, and business-level models are built only on cleansed data. It works because it makes every table's rebuild path obvious.

Lake, warehouse, lakehouse describe where the table lives. A warehouse owns its storage and gives you transactions, a planner and governance in one system. A lake is files in object storage that anything can read, with correspondingly weak guarantees. The lakehouse pattern puts a table format over lake files to recover warehouse-like guarantees while keeping open storage. Each is a different answer to one question: how much do you value engine independence versus integration?

Ownership: central platform or data mesh

A centralized team owns every pipeline: consistent, easy to govern, and a bottleneck that grows with the number of source systems, because the central team is always the least informed party about any given source's semantics. The decentralized alternative (commonly called a data mesh) pushes ownership to the teams that produce the data, and makes the platform team responsible for the substrate and the standards rather than the content.

Decentralization only works when three things are real: published data contracts with explicit schemas and compatibility rules, a shared platform so each team is not inventing its own stack, and a catalog so consumers can find what exists. Without those, a mesh is just an unowned pile with better branding.

The tradeoffs that actually decide the design

  • Freshness against cost and complexity. Every step toward real time removes a batching opportunity, and batching is what makes analytics affordable. Ask what decision is made with the data and how often; a dashboard someone reads each morning does not need a streaming pipeline.
  • Schema-on-write against schema-on-read. Enforcing schema at ingest rejects bad data early and blocks the producer. Accepting anything and interpreting later keeps ingestion resilient and moves the failure into a query that someone will debug later, usually under pressure.
  • Copies against coupling. One canonical table means every consumer is coupled to it. Per-consumer copies decouple teams and then quietly disagree. Neither is free; pick deliberately and write down which one you picked.
  • Open formats against integrated systems. Open formats preserve the option to change engines. Integrated systems remove whole categories of operational work. This is a reversibility decision, not a performance one.

Common failure modes

  • The small-file problem. A streaming writer commits every few seconds and creates millions of tiny objects; reads degrade steadily until a scheduled compaction exists.
  • Partitioning by a column nobody filters on. All the maintenance cost, none of the pruning benefit.
  • Silent schema drift. An upstream field changes type or meaning, the loader widens the column or nulls it, and no alert fires because the row count is unchanged. Count-based monitoring cannot see this; schema and distribution checks can.
  • No backfill story. The pipeline can run forward but cannot recompute the past, so every logic fix creates a permanent discontinuity in the data.
  • The orchestrator used as the compute engine. Pulling rows into the scheduler process to transform them turns a coordination tool into a single-node bottleneck with no recovery semantics.
  • A "single source of truth" that is five copies. Usually visible when two dashboards disagree and nobody can say which table is canonical.
  • Unbounded raw retention with no lifecycle. Storage is cheap per unit and expensive in aggregate; the fix is a retention and tiering policy written at design time, not after the invoice.

When you do not need any of this

Most organizations that think they have a big data problem have a query problem. If the working set fits on one machine, a single well-indexed relational database with a read replica for analytics will outperform a distributed platform on every axis that matters — latency, correctness, operational burden and the number of people required to keep it alive. The honest triggers for distributed architecture are: the data genuinely does not fit or cannot be processed in the available window; independent teams need to publish and consume data without coordinating releases; or regulatory requirements demand retention and lineage that the transactional system cannot provide. Growth alone is not a trigger. Reach for a warehouse before a lake, and a scheduled pipeline before a streaming platform.