AWS Lesson 104 of 123

AWS Enterprise Architecture: Big Data Processing

In a nutshell

Imagine a professional kitchen that owns one enormous walk-in fridge but rents its ovens by the minute. The fridge — your S3 data lake — holds every ingredient permanently and never gets thrown out. The ovens — your Spark compute — are wheeled in only when there is a dish to cook, run hot for exactly as long as the cooking takes, and are wheeled straight back out the moment the timer dings, so you stop paying for them. The old way was to buy a hundred ovens, keep them all switched on 24/7 in case a banquet showed up, and pay the electricity bill whether you cooked or not. This lesson builds the rent-by-the-minute kitchen: permanent storage that never moves, and disposable compute that is born for a single job and dies when it finishes.

“Big data processing” just means you have more data than one computer can chew through in the time you have, so you split the work across many machines — that is what Apache Spark does — and you rent those machines only for the job. On AWS the fridge is Amazon S3; the ovens are Amazon EMR (Spark on rented EC2 servers), EMR Serverless (Spark with no servers to think about), and AWS Glue (managed Spark for routine pipelines). A shared recipe index — the Glue Data Catalog — tells every oven what tables exist, and a shared lock — AWS Lake Formation — decides who may read which columns and rows, even from inside a Spark job. Amazon Athena lets you run SQL straight over the fridge with no oven at all.

The one idea to carry through the whole lesson: storage is permanent and compute is disposable. If a beginner remembers only that, everything else — Spot pricing, data skew, medallion zones, governance — is detail about how you make the disposable half cheap and reliable.

Level: Advanced · Time: ~54 min

Before this lesson, it helps to know: what Amazon S3 is (objects, buckets, prefixes), the basic idea of an IAM role, and that SQL queries read tables. If “distributed processing” is new to you, that is fine — this lesson explains it. The companion Data Lakehouse on AWS lesson covers the storage side (open table formats, store-once-query-many) in depth, and the Real-Time Streaming lesson covers sub-second event handling; this one is the batch processing / compute half.

After this lesson you will be able to:

Big data processing fails in a way that has almost nothing to do with the data and everything to do with the compute. A team stands up a long-running EMR cluster sized for the year-end reprocessing job, runs the nightly 40-minute pipeline on it, and then pays for 240 idle CPU-hours a day. Or the opposite: they pick a cluster that was right for the demo, the data grows 3x, one Spark stage spills to disk because a single skewed key landed on one executor, and a job that should take 20 minutes runs for six hours and then OOM-kills itself at 95%. Neither problem is a Spark problem. They are the absence of a processing architecture — a deliberate separation between durable storage that never moves, ephemeral compute that you turn on for the duration of a job and tear down after, a serverless tier for the spiky and the small, and a catalog-plus-governance plane so every engine sees the same tables and the same access rules. This article builds that platform end to end on Amazon EMR, S3, AWS Glue, Amazon Athena, AWS Lake Formation, and Apache Spark, and treats the hard parts — cluster topology, Spot economics, data skew, job orchestration, and “which of the three EMRs do I actually run” — as first-class concerns.

The single most important idea is this: storage is permanent and compute is disposable. Your data lives on S3 forever; the Spark cluster that transforms it should exist only while a job runs and cost nothing the rest of the day. Get that decoupling right and big-data processing becomes an exercise in picking the cheapest, correctly-sized, transient compute for each job — and the multi-million-dollar Hadoop cluster that ran 24/7 whether or not anyone had work for it becomes a line item you no longer pay.

This is the processing-and-compute companion to the AWS Data Lakehouse reference (which centres on storage, open table formats, and store-once-query-many) and the Real-Time Streaming reference (which centres on sub-second event handling). Here the spine is the batch Spark pipeline and the question is how the heavy transformation actually runs.

The business scenario

Picture an organisation that has more data than its current tools can chew through in the time it has. The shape is identical at 40 engineers and at 4,000.

The early-stage version. A 60-person ad-tech startup lands ~4 TB/day of impression and click logs into S3. They began with a single r5.4xlarge EC2 box running a Python script and a cron job. It worked at 200 GB/day. At 4 TB/day the nightly aggregation now takes 9 hours, finishes after the morning standup, and falls over roughly twice a week when a partner sends a malformed batch. Every attempt to “just make the box bigger” buys a few weeks before the data outgrows it again. The founding engineer knows the answer is “distributed processing” but has watched friends drown in a self-managed Hadoop cluster and does not want to hire a platform team to babysit YARN.

The large-enterprise version. A telco with 60M subscribers runs a 400-node on-premises Hadoop cluster that a team built in 2017. It processes call-detail records, network telemetry, and billing events — roughly 600 TB/month landing, 12 PB at rest. The cluster runs at 30% average utilisation because it is sized for the monthly fraud-detection reprocessing peak; the other 29 days it is mostly idle but fully powered and fully staffed. Hardware refresh is a capital project with an 18-month lead time. Three different teams want more capacity and there is none to give because adding nodes means buying racks. The CFO sees a fixed, enormous, growing cost and the data team sees a queue of jobs they can’t schedule.

Both organisations have the same structural problem, and — exactly as with the lakehouse — it is a coupling problem, but on a different axis:

The processing platform in this article decouples all four. One durable copy of data on S3. Compute that is born for a job and dies when it finishes — a transient EMR cluster, an EMR Serverless application that scales to zero, or a Glue serverless job. Each workload on the compute shape sized for it. And per-job isolation so a poison batch in the fraud pipeline can’t touch the billing pipeline. You stop paying for peak capacity at idle, scaling becomes an API call, and the heaviest job in the year no longer dictates what you pay in July.

The promise to the business: transform any volume of data, on compute you turn on for exactly as long as the job runs, sized per job, and pay by the second of actual work — not by the calendar.

Architecture overview

The platform is a durable S3 data lake in the centre, surrounded by disposable compute engines that all read and write the same data through one shared catalog and one governance plane. Read it as five tiers, with two cross-cutting planes (catalog/governance underneath, orchestration/observability on top) spanning the whole width.

AWS big-data processing platform: source systems and ingestion (DMS, Firehose, Glue connectors) land data in an S3 raw/processed/curated spine; a disposable Spark compute fleet (EMR on EC2 with Spot, EMR Serverless, AWS Glue, EMR on EKS) reads and writes the same tables through a Glue Data Catalog and Lake Formation governance plane; Athena, Redshift and SageMaker serve the curated edge, with Step Functions/MWAA orchestration and an IAM/KMS/VPC/CloudTrail baseline.

Tier 1 — Ingestion into the raw lake (S3). Source data lands, untransformed, in a raw zone on S3. Bulk and relational sources arrive via AWS DMS (CDC from operational databases) or direct S3 uploads/partner drops; high-volume event data arrives through Kinesis Data Firehose writing partitioned objects; SaaS sources come through AWS Glue connectors or AppFlow. Nothing is processed here. Raw is immutable, append-only, partitioned by ingest date, and is the system of record for replay — every derived dataset can be rebuilt from it.

Tier 2 — The lake on S3 (the spine). S3 holds the data across curated zones — raw (as-landed), processed/cleaned (conformed, deduplicated, columnar), and curated/aggregated (business-level marts and feature tables). Curated data is written as Parquet (columnar, compressed with zstd/Snappy) and, where transactional semantics matter, as Apache Iceberg tables. This is the one physical copy of all data; every compute engine is a temporary lens over it. (The storage design — zone layout, table formats, small-file economics — is the subject of the lakehouse reference; here we focus on the engines that read and write it.)

Tier 3 — The processing fleet: where the heavy lifting actually happens. This is the centre of gravity of this architecture and where it differs from every other data reference. Three Spark-capable compute shapes read the same S3 + catalog and write back the same tables, each chosen per workload:

A fourth shape, EMR on EKS, runs Spark as pods on a shared Amazon EKS cluster — the right answer when an org has standardised on Kubernetes and wants Spark to share capacity, tooling, and Karpenter-driven autoscaling with its other containerised workloads. All four write the same Parquet/Iceberg tables to the same Glue catalog; the engine is a per-workload choice, not a platform-wide religion.

Tier 4 — Catalog and governance (the plane underneath everything). The AWS Glue Data Catalog is the single technical metadata store — every table, schema, and partition registered once, read by EMR, Glue, Athena, and Redshift alike. AWS Lake Formation sits on top as the permission layer: it owns the S3 locations and grants database/table/column/row-level access rather than raw S3 access. Every engine — including an EMR cluster running Spark — vends credentials through Lake Formation at runtime, so “analysts cannot see the msisdn column; the EU team sees only EU rows” is enforced identically across Spark, Glue, and Athena.

Tier 5 — Serving and consumption (the right edge). Processed and curated tables are served to three patterns: Amazon Athena (serverless, pay-per-TB-scanned SQL for ad-hoc exploration and data-science queries — zero standing cost), Amazon Redshift / Spectrum (warehouse-grade BI over curated marts joined to hot dimensions), and Amazon SageMaker / EMR (reading curated feature tables directly for ML training). BI tools — QuickSight, Tableau, Power BI — sit beyond the engines. The data is never permanently copied into any engine; S3 is the one copy.

The two cross-cutting planes. Orchestration sits on top: AWS Step Functions and Amazon Managed Workflows for Apache Airflow (MWAA) sequence the pipeline — land raw, run Glue/EMR transforms, validate data quality, publish curated, then signal downstream — with retries, dependencies, and backfills. Observability spans everything: EMR/Glue push Spark metrics and logs to CloudWatch, the Spark History Server and EMR’s managed UIs expose stage-level execution, and Glue Data Quality gates the pipeline before bad data reaches consumers.

The diagram, in words. Picture a wide canvas. On the left, a column of source systems (databases, partner feeds, event streams, SaaS) with arrows through DMS / Firehose / Glue connectors into a single large S3 “raw” bucket. In the centre, S3 is drawn as the spine — three stacked buckets labelled raw → processed → curated, the unmistakable middle of the picture. Sitting above and around the S3 spine, four compute boxes all draw bidirectional arrows to S3: EMR on EC2 (drawn with a small On-Demand core block and a larger dashed “Spot task nodes” block to signal transience), EMR Serverless (a cloud with no fixed shape), AWS Glue (a gear), and EMR on EKS (Spark pods inside an EKS hexagon). Crucially, each compute box is drawn with a dashed border to convey ephemeral — they appear for a job and vanish. Underneath the entire spine runs a full-width band: Glue Data Catalog (one source of schema) with Lake Formation layered on it (the lock), and every compute box’s arrow to S3 passes through this band. On the right, Athena / Redshift / SageMaker read curated S3, feeding QuickSight / BI at the far edge. Across the top, an orchestration ribbon — Step Functions / MWAA (Airflow) — with control arrows reaching down into each compute box, and a parallel CloudWatch / Spark History / Glue Data Quality observability ribbon. Cross-cutting boxes at the very bottom — IAM Identity Center, KMS, VPC, CloudTrail — touch every tier. The defining visual: a permanent storage spine, ringed by disposable compute that is born and dies per job, all governed through one plane.

The ingestion menu in depth

Tier 1 named the three most common on-ramps — DMS, Firehose, and Glue connectors — but a real platform meets data wherever it lives, and “how does it get into the raw zone” is a decision you make per source, not once. Here is the fuller menu and when each is the right door.

Source shape Service Why it fits
Operational database, ongoing change capture AWS DMS (CDC) Streams inserts/updates/deletes from Oracle/PostgreSQL/MySQL/SQL Server into S3 as they commit; full-load + ongoing replication
High-volume event/log streams Amazon Data Firehose (formerly Kinesis Data Firehose) Buffers and writes partitioned objects to S3 with no code; can convert JSON→Parquet in flight and now land directly into Iceberg tables
SaaS apps (Salesforce, ServiceNow, Zendesk) Amazon AppFlow / Glue connectors Managed, schema-aware pulls with no API client to write
On-prem file servers / HDFS / other object stores, online AWS DataSync Managed, accelerated, scheduled network copy over the internet or Direct Connect; handles NFS/SMB/HDFS/S3-compatible
Inbound files from external partners over SFTP/AS2 AWS Transfer Family Managed SFTP/FTPS/FTP/AS2 endpoints that land straight into S3 — no SFTP servers to run
Petabytes where the network would take weeks AWS Snow Family (Snowball Edge) Ship a physical, encrypted appliance; the wire is too slow or too costly for the first bulk load

The decision that catches teams out is online vs. offline for the first big load. DataSync moves data over the network; Snowball ships a box. The rule of thumb is to estimate the transfer time first: 100 TB over a fully saturated 1 Gbps link is roughly 9 days of continuous transfer at ideal throughput, and real links rarely sustain ideal. If the initial copy would take more than a week or two — or if pushing it over the wire would blow your Direct Connect / internet-egress budget — ship a Snowball Edge for the historical backfill and use DataSync (or DMS CDC) for the ongoing delta. Meridian’s 12 PB history in the reference example is squarely Snow-family territory for the seed; the ~20 TB/day that lands afterward is DMS + Firehose.

Worked example — wiring one source correctly. The billing system is Oracle and must be captured continuously without a nightly full export hammering the OLTP database. The right pattern: an AWS DMS task in full-load + CDC mode reads the Oracle redo logs, does one initial full copy into raw/billing/, then streams every subsequent change as it commits. DMS writes CDC output to S3 as timestamped files partitioned by ingest date. A downstream Glue job then merges those change files into an Iceberg table in the processed zone with MERGE INTO, so updates and deletes are applied, not merely appended — turning an append-only change stream into a correct current-state table. The key insight: ingestion lands change events; the processing tier turns them into tables. Trying to make the ingestion layer emit perfectly-merged current state is where brittle pipelines come from.

Everything lands in the raw zone untransformed — and that immutability is exactly what lets you replay. If a merge-logic bug corrupts the processed billing table six months from now, the fix is to re-run the merge from raw, not to restore from backup, because raw is the system of record and every derived table is recomputable from it.

Component breakdown

Component Role in the platform Key configuration choices
Amazon S3 (zoned) The permanent spine — one physical copy of raw/processed/curated data; every engine’s read/write target Separate buckets/prefixes per zone; partition by date; S3 Intelligent-Tiering on raw; lifecycle to Glacier for cold; Block Public Access + SSE-KMS; columnar Parquet for processed/curated
EMR on EC2 (transient) Heavy/scheduled Spark batch — big ETL, reprocessing, ML feature jobs; launched per job, terminated after Instance fleets with diversified Spot for task nodes; On-Demand for primary/core; managed scaling; --auto-terminate; EMRFS + Glue Catalog; runtime roles for Lake Formation
EMR Serverless Spark with no cluster to manage — spiky, intermittent, small-to-medium jobs; scales to zero Pre-initialised capacity for latency-sensitive jobs vs. cold-start for cheap; per-job application; vCPU/GB-second billing; auto-stop idle apps
AWS Glue ETL Serverless Spark for routine raw→processed→curated transforms and orchestration glue Glue 5.0 (Spark 3.5); job bookmarks for incremental; Auto Scaling workers; Flex execution for non-urgent jobs (cheaper); Glue connectors for sources
EMR on EKS Spark-as-pods on a shared EKS cluster — for K8s-standardised orgs Virtual clusters per team; Karpenter + Spot node pools; pod templates; shares tooling/observability with other containerised workloads
Apache Spark The processing engine inside EMR/Glue/EKS — the distributed compute itself Tune partitions/shuffle; Adaptive Query Execution (AQE) for skew + dynamic coalescing; broadcast joins for small dims; cache hot intermediates; EMR’s optimised Spark runtime
Glue Data Catalog One technical metadata store every engine shares One catalog per environment; databases per domain; crawlers for schema discovery on raw; tables registered natively for curated
Lake Formation Governance plane — fine-grained, engine-agnostic permissions enforced even inside Spark Register S3 locations; remove IAMAllowedPrincipals; LF-Tags for scale; column masking + row filters; EMR runtime role enforcement
Amazon Athena Serverless ad-hoc SQL over processed/curated — the default exploration engine Engine v3 (Trino); workgroups with per-query scan caps; results to dedicated S3; query-result reuse cache
Step Functions / MWAA Orchestration — sequence, retry, branch, backfill the pipeline Step Functions for AWS-native, event-driven flows; MWAA (Airflow) for complex DAGs/cross-system dependencies; both trigger EMR/Glue and gate on data quality
Glue Data Quality Data-quality gates between zones, before bad data reaches consumers DQDL rules on processed→curated; fail-fast on freshness/null-rate/referential checks; quality results to CloudWatch/EventBridge

Why each is here, and the choices that matter:

S3 is the spine, and it never moves. Every other component in the diagram is disposable; S3 is the one thing that persists. Because raw is immutable and replayable, and processed/curated are recomputable from raw, S3 is simultaneously your storage, your backup substrate, and your DR foundation. The single most consequential storage choice for processing performance is columnar Parquet with compression and sensible partitioning: a Spark job that reads partition-pruned, predicate-pushed Parquet touches a fraction of the bytes of one reading raw JSON, which is the difference between a 20-minute job and a 3-hour one.

The three-plus-one compute shapes are the heart of the architecture — and choosing between them is the skill. This table is the decision most teams get wrong:

If the workload is… Run it on… Because…
Large, scheduled, predictable (nightly multi-TB ETL); long reprocessing EMR on EC2, transient + Spot Best $/compute via Spot task nodes and instance-fleet diversification; cluster lives only for the job
Spiky, intermittent, unpredictable, or you don’t want to size instances EMR Serverless Scales to zero between runs; no capacity planning; pay per vCPU/GB-second of actual work
Routine declarative ETL on a schedule/event; incremental loads AWS Glue Job bookmarks, connectors, fastest to author; serverless; Flex for cheap non-urgent runs
Ad-hoc SQL exploration; “what’s in this table?” Athena Zero standing cost; pay per TB scanned; no Spark to manage at all
Your org runs everything on Kubernetes already EMR on EKS Spark shares capacity, Karpenter autoscaling, and tooling with other K8s workloads

The anti-pattern is choosing one shape for everything: running ad-hoc exploration on a giant standing EMR cluster (you’ll pay for idle), or forcing a steady nightly ETL onto EMR Serverless when a transient Spot cluster would be 70% cheaper, or hand-managing EMR for a job Glue would author in an afternoon.

Transient EMR + instance fleets + Spot is where the money is saved. A transient cluster launches, runs, and self-terminates, so there is no idle cost. Instance fleets let you specify a target capacity and a menu of acceptable instance types across multiple AZs; EMR fills it from whatever Spot capacity is cheapest and available, dramatically reducing the chance that a Spot shortage stalls your job. The rule that keeps this reliable: primary and core nodes on On-Demand (they hold HDFS/shuffle data and the application master — losing them kills the job); task nodes on Spot (they only run compute — losing one just reruns its tasks). Get that split right and you capture 60–90% Spot savings on the bulk of the cluster without risking job completion.

Spark is the engine, and AQE is the feature that earns its keep. The classic big-data failure — one stage runs for hours while 199 of 200 tasks finished in seconds — is data skew: a single hot key (one giant advertiser, one heavy subscriber) lands all its rows on one executor. Spark’s Adaptive Query Execution detects skewed partitions at runtime and splits them, dynamically coalesces tiny shuffle partitions, and switches join strategies based on actual sizes — turning a class of six-hour-then-OOM jobs into ones that just finish. Combined with broadcast joins for small dimensions and right-sized shuffle partitions, AQE is the single highest-leverage Spark setting on this platform.

Glue Catalog is the contract; Lake Formation is the lock — and the subtlety is that it works inside Spark. On a self-managed Hadoop cluster, once a job has the cluster’s IAM role it can read any S3 it can reach — governance is coarse and per-cluster. Here, an EMR cluster runs with a runtime role, and Lake Formation vends scoped, temporary credentials per query even to Spark, so the same column mask that hides msisdn in Athena also hides it from a Spark SELECT on EMR. Policy lives in one place and is enforced everywhere, including the heavy-compute tier.

Orchestration is not optional at enterprise scale. A real pipeline is a DAG: land raw, dedupe, conform, join reference data, aggregate, validate, publish, notify. Step Functions expresses this as an AWS-native state machine (great for event-driven, EMR/Glue-centric flows); MWAA (managed Airflow) is the choice when you have complex cross-system dependencies, hundreds of DAGs, or a team already fluent in Airflow. Both give you retries, backfills, and the ability to rerun one failed step rather than the whole pipeline — and both gate on Glue Data Quality so a failed freshness or null-rate check stops the pipeline before bad data reaches a dashboard.

Open table formats and the medallion lake, in depth

The overview called the zones raw → processed → curated. That naming is one dialect of a pattern the industry more often calls medallion: bronze (raw, as-landed), silver (cleaned, conformed, deduplicated), and gold (curated business-level marts and features). Same idea — and it matters for a reason beyond tidiness: each promotion between zones is a quality and trust boundary. Bronze is faithful to the source, warts and all; silver is where schema is enforced and duplicates die; gold is what a dashboard or a model is allowed to touch. When someone asks “can I trust this number,” the honest answer is “it depends which zone it came from,” and the zone tells you.

Why plain Parquet files are not enough, and what an open table format adds. A directory of Parquet files has no notion of a transaction. Two jobs writing the same partition at once can leave a reader seeing half of each. There is no atomic “replace yesterday’s data,” no way to UPDATE or DELETE a single row for a GDPR erasure request without rewriting whole partitions by hand, and no reliable snapshot to query “as of last Tuesday.” An open table format — Iceberg, Hudi, or Delta Lake — adds a metadata layer on top of the Parquet files that gives you ACID transactions, row-level MERGE/UPDATE/DELETE, schema evolution, hidden partitioning, and time travel, while the data itself stays as open Parquet on S3 that any engine can read.

Format Origin & sweet spot Notable strengths
Apache Iceberg Netflix; the AWS-default open format Hidden partitioning + partition evolution (change partitioning without rewriting queries), snapshot isolation, broadest AWS support (Athena, EMR, Glue, Redshift Spectrum, S3 Tables)
Apache Hudi Uber; upsert- and streaming-heavy ingestion Record-level indexing, copy-on-write vs. merge-on-read tradeoff, incremental pulls; strong when you ingest a firehose of updates
Delta Lake Databricks; ubiquitous where Databricks runs Transaction log (_delta_log), mature tooling, OPTIMIZE/Z-ORDER; first-class on Databricks, readable on EMR

For a green-field AWS lake, Iceberg is the safe default — it has the deepest native integration across Athena, EMR, Glue, and Redshift Spectrum, and it is the format AWS chose for S3 Tables and SageMaker Lakehouse. Choose Hudi when your dominant pattern is high-frequency upserts / streaming CDC and you want record-level indexing; choose Delta when you are already a Databricks shop. The Deploy Apache Iceberg on S3 with Glue Catalog and compaction lesson works the Iceberg mechanics end to end.

S3 Tables — managed Iceberg storage. Announced at re:Invent 2024, Amazon S3 Tables is a purpose-built S3 bucket type (a “table bucket”) that stores Iceberg tables and runs the housekeeping for you: automatic compaction of small files, snapshot expiration, and unreferenced-file cleanup, with tighter integration into the Glue Data Catalog and Lake Formation. The pitch is that the two chores that wreck self-managed Iceberg lakes — small-file compaction and snapshot/garbage cleanup — become a managed background service, at the cost of some control and a storage-management fee. On a general-purpose bucket you own those maintenance jobs yourself.

The small-files problem, with the numbers. This is the single most common self-inflicted wound on a lake, so it is worth making concrete. Suppose a streaming job writes one 4 MB Parquet file every 10 seconds into a partition. After a day that partition holds ~8,600 files totalling ~34 GB. A Spark job that reads it must list every object (S3 LIST is paginated and not free), open every file (each open is a round-trip, and every Parquet file carries footer-parsing overhead), and it gets almost none of Parquet’s columnar benefit because each file is tiny. The same 34 GB written as ~34 files of ~1 GB reads dramatically faster — fewer LISTs, far fewer opens, bigger contiguous columnar reads, better predicate push-down. The fix is compaction: a scheduled job (or S3 Tables’ automatic compaction, or Iceberg’s rewrite_data_files) that rewrites many small files into few right-sized ones. The target from the implementation section — 128 MB–1 GB per file — is not a nicety; it is the difference between a query that scans efficiently and one that dies of a thousand paper cuts. Throwing a bigger cluster at a small-files problem does nothing, because the bottleneck is object count and metadata, not CPU.

Implementation guidance

Bucket and zone layout. Consistent prefixing across zones, partitioned for pruning:

s3://acme-bd-raw-<acct>/<source>/<table>/ingest_date=YYYY-MM-DD/
s3://acme-bd-processed-<acct>/<domain>/<table>/dt=YYYY-MM-DD/   # Parquet
s3://acme-bd-curated-<acct>/<domain>/<table>/                  # Parquet/Iceberg

Raw partitions by ingest date for cheap lifecycle and replay; processed/curated partition by the business date queries actually filter on. The processing rule of thumb: aim for 128 MB–1 GB Parquet files — too small and Spark wastes time on task overhead and S3 listing; too large and you lose parallelism.

Terraform is the right IaC here (the team standardises on it). Manage as code:

Transient EMR cluster spec — the choices that matter (expressed conceptually; this lives in the orchestration step, not standing infra):

{
  "ReleaseLabel": "emr-7.5.0",                 // Spark 3.5, optimised runtime
  "Applications": ["Spark"],
  "AutoTerminationPolicy": { "IdleTimeout": 600 },  // die if idle 10 min
  "ManagedScalingPolicy": { "MinCapacity": 2, "MaxCapacity": 64 },
  "Instances": {
    "InstanceFleets": [
      { "InstanceFleetType": "MASTER", "TargetOnDemandCapacity": 1,
        "InstanceTypeConfigs": [ { "InstanceType": "m6g.xlarge" } ] },
      { "InstanceFleetType": "CORE",   "TargetOnDemandCapacity": 2,   // On-Demand: holds data
        "InstanceTypeConfigs": [ { "InstanceType": "r6g.2xlarge" } ] },
      { "InstanceFleetType": "TASK",   "TargetSpotCapacity": 40,      // Spot: pure compute
        "InstanceTypeConfigs": [                                       // diversified menu
          { "InstanceType": "r6g.4xlarge" }, { "InstanceType": "r5.4xlarge" },
          { "InstanceType": "r6i.4xlarge" }, { "InstanceType": "m6g.8xlarge" }
        ],
        "LaunchSpecifications": { "SpotSpecification": {
          "AllocationStrategy": "price-capacity-optimized" } } }   // best Spot strategy
    ]
  },
  "Configurations": [ { "Classification": "spark-defaults", "Properties": {
    "spark.sql.adaptive.enabled": "true",
    "spark.sql.adaptive.skewJoin.enabled": "true",            // the skew killer
    "spark.sql.adaptive.coalescePartitions.enabled": "true",
    "spark.dynamicAllocation.enabled": "true"
  } } ]
}

The headline decisions: price-capacity-optimized Spot allocation (balances cheapest with least-likely-to-be-interrupted), On-Demand core / Spot task split, managed scaling so the cluster grows to the data and shrinks after, auto-termination so it never idles, and AQE on so skew and small-partition problems resolve at runtime.

EMR Serverless removes all of the above for spiky jobs — you submit a Spark job to an application and AWS handles workers:

aws emr-serverless start-job-run \
  --application-id <app> --execution-role-arn <runtime-role> \
  --job-driver '{ "sparkSubmit": {
      "entryPoint": "s3://acme-bd-code/jobs/daily_rollup.py",
      "sparkSubmitParameters": "--conf spark.sql.adaptive.enabled=true" } }'

Use pre-initialised capacity when a job is latency-sensitive (workers stay warm); accept cold-start when it’s a cheap background job.

Networking. Keep all data-plane traffic on the AWS backbone via a VPC Gateway Endpoint for S3 (free) and Interface (PrivateLink) endpoints for Glue, Lake Formation, Athena, KMS, STS, and the EMR APIs. EMR clusters, EMR Serverless, and EMR-on-EKS nodes run in private subnets with no public IPs; a NAT gateway handles only OS/package egress. The result: terabytes of shuffle and S3 I/O never traverse the public internet, and you avoid per-GB internet data charges on the heaviest traffic in the system.

Identity wiring. Human access flows from the IdP (the org standardises on Entra ID) through IAM Identity Center (SAML/SCIM) → permission sets → data-access roles, which are then granted Lake Formation permissions (not S3 permissions). Pipeline identities are scoped per engine: each Glue job, each EMR cluster (via its runtime role), and each EMR Serverless application gets its own role with only the Lake Formation grants it needs. No human or job role gets direct s3:GetObject on the lake buckets — Lake Formation vends scoped, temporary credentials per query, even into Spark. This is the wiring decision that makes governance real rather than theatrical, and it’s the one that most distinguishes this from a legacy Hadoop cluster where the cluster role could read everything.

A worked example: reading a Spark job, and what it costs

Two things separate people who use this platform from people who merely provision it: the ability to read why a Spark job is slow, and the ability to put a dollar figure on the compute choice. Here is one of each, worked end to end.

Reading a skewed join, step by step

The scenario: nightly, you join a 2 TB impressions table to a 500 MB advertisers dimension, group by advertiser, and write daily rollups. It ran in 40 minutes for months, then one week it takes 5 hours and sometimes OOM-kills at 95%. Nothing about the cluster changed. Walk it through:

  1. The join. advertisers is 500 MB — larger than Spark’s default broadcast threshold (spark.sql.autoBroadcastJoinThreshold, default 10 MB), so Spark plans a shuffle (sort-merge) join: both sides are hash-partitioned by advertiser_id across (by default) 200 shuffle partitions, and matching keys are shipped to the same partition. So far, normal.
  2. The skew. One advertiser — a giant retail brand — accounts for 22% of all impressions. Every one of those rows hashes to the same shuffle partition. So 199 partitions get ~9 GB each and finish in a minute, while one partition gets ~450 GB and lands on a single executor. That one task runs for hours; the stage cannot finish until it does; and if the executor’s memory + local disk cannot hold ~450 GB of shuffle, it spills and eventually OOMs. This is data skew, and “199 of 200 tasks done in seconds, one running for hours” is the signature you learn to recognise in the Spark UI.
  3. The fix — AQE skew join. With spark.sql.adaptive.enabled=true and spark.sql.adaptive.skewJoin.enabled=true, Spark measures partition sizes at runtime and, when a partition exceeds both skewedPartitionThresholdInBytes (default 256 MB) and skewedPartitionFactor × the median partition size (default factor 5), it splits that partition into several and replicates the matching dimension rows to each — so the hot advertiser is processed by many tasks in parallel instead of one. The 5-hour job returns to ~40 minutes without you touching a line of business logic.
  4. The other lever — broadcast. The 500 MB dimension only needed a shuffle because it was “big.” If you project it to just the columns and rows actually used (say 40 MB) and raise the broadcast threshold, Spark can broadcast the dimension to every executor and skip the shuffle entirely — no skew is possible, because the big side is never shuffled. Broadcasting is the cheapest join when one side is small; the art is knowing when a side is small enough.

The lesson: skew is a data-distribution problem, not a capacity problem. A bigger cluster gives the hot partition a bigger single executor, buying a little headroom before OOM, but it never parallelises that one task — only splitting the partition (AQE), changing the join (broadcast), or pre-salting the key does. That is why “the job is slow, add nodes” is the most expensive way to not fix a skewed job.

Costing the compute choice

Now the money. The same nightly job, four ways, on a representative ~40-node r6g.4xlarge-class footprint that runs ~1 hour a night. Figures are illustrative — regions and rates vary, so check current pricing — but the ratios are the point, and they hold everywhere. Assume ~$1 per node-hour On-Demand (EC2 + the EMR uplift) and Spot at ~30% of that.

Option What you pay for Rough monthly cost Notes
Standing cluster 24/7 40 nodes × 24 h × 30 d ~$29,000 Idle 23 h/day; the legacy default
Transient EMR, On-Demand only 40 nodes × 1 h × 30 d ~$1,200 No idle — the structural win
Transient EMR, Spot task nodes ~4 On-Demand core + ~36 Spot task ~$350–500 60–90% off the compute bulk; the production default
EMR Serverless vCPU-second + GB-second of the job ~$500–700 No sizing at all; slightly dearer than tuned Spot for a steady large job

Two conclusions fall out. First, the giant saving is transience, not tuning: moving from a standing cluster to any transient option removes ~95% of the bill by deleting idle — that is the ~$29,000 → ~$1,200 step, and it is why the whole architecture exists. Second, among the transient options the choice is economic, not moral — Spot on a right-sized transient cluster is cheapest for this steady, predictable, large job, which is exactly why Meridian ran its nightly job on transient EMR + Spot and reserved Serverless for the unpredictable analyst jobs where “scale to zero, no sizing” beats Spot’s slightly lower rate. Add the Athena serving path — $5 per TB scanned (representative), cut ~90%+ by Parquet + partition pruning + compression — and the platform’s entire cost story is three multiplications you can do on a napkin.

Enterprise considerations

Security and Zero Trust. The model is “no implicit S3 access; every read — including from Spark — is brokered.” Concretely: (1) Block Public Access on every bucket, account-wide; (2) SSE-KMS with separate CMKs per zone for independent revocation and per-zone CloudTrail audit; (3) Lake Formation as the only path to data, so even an EMR cluster’s runtime role holds Lake Formation grants, not bucket policies; (4) column masking on PII (msisdn, imsi, card_pan) and row filters for tenancy/region, enforced in Spark, Glue, and Athena alike; (5) all traffic on PrivateLink; (6) EMR security configurations add in-transit and at-rest encryption for the cluster’s own intermediate data (local disks, shuffle); (7) CloudTrail data events + Lake Formation’s audit log give a complete “who processed/read which column when” trail. LF-Tags keep this scalable — tag a table confidentiality=pii once and grant on the tag, so new tables inherit policy automatically.

Cost optimization — this is the platform’s reason to exist, and it’s layered.

Scalability. S3 absorbs effectively unlimited objects and throughput; partition-prefix design avoids hotspots. EMR managed scaling grows a cluster mid-job as a stage demands more executors and shrinks it after; EMR Serverless and Glue absorb concurrency automatically; Athena is serverless. The real scaling discipline is not capacity — it’s data engineering: keeping Parquet file sizes in the sweet spot (compaction jobs to fix small files), partitioning on the columns queries filter, and handling skew with AQE so adding nodes actually helps instead of one hot key bottlenecking on one executor regardless of cluster size. Throwing hardware at a skewed job is the most expensive way to not fix it.

Reliability and DR (RTO/RPO).

Observability. EMR and Glue emit Spark metrics and driver/executor logs to CloudWatch; the Spark History Server (and EMR’s managed Spark UI / Persistent UIs) expose stage-level execution — the single most important diagnostic for “why is this job slow,” surfacing spills, skew, and shuffle volume. Glue Data Quality (DQDL) gates the processed→curated boundary so freshness/null-rate/referential failures stop the pipeline and fire EventBridge alerts before bad data reaches dashboards. Track four SLOs: pipeline freshness (curated lag), job cost (compute-hours and bytes scanned/day), data-quality pass rate, and Spot interruption rate (rising interruptions are a signal to broaden the instance-fleet menu).

Governance. Lake Formation LF-Tags are the scalable model — a taxonomy (domain, confidentiality, retention) applied to databases/tables, grants written against tags so policy is inherited not hand-maintained. This is also the foundation for a data-mesh evolution: each domain owns its processed/curated databases and the jobs that produce them, sharing across domains via Lake Formation cross-account grants, with a central platform team owning only the tag taxonomy, the orchestration substrate, and the EMR/Glue golden patterns. Pair with Amazon DataZone for business-level discovery and access requests on top of the technical Glue catalog.

Reference enterprise example

Meridian Telecom is a fictional mobile carrier with 62M subscribers, running the legacy mess described at the top: a 400-node on-prem Hadoop cluster built in 2017, ~600 TB/month landing (call-detail records, RAN telemetry, billing events, app analytics), 12 PB at rest, ~30% average utilisation because it’s sized for the monthly fraud-reprocessing peak. The data-platform run-rate (hardware amortisation, datacentre, ops staff, Hadoop support) is ~$6.4M/year, three teams are queued for capacity that doesn’t exist, and the last audit flagged that the cluster role could read subscriber imsi/msisdn with no fine-grained control.

What they built. Over three quarters they migrated to the architecture above:

Decisions worth noting. They put core nodes on On-Demand and task nodes on Spot after an early pilot where an all-Spot cluster lost its core nodes mid-job and had to restart from scratch. They standardised on Graviton (m6g/r6g) for ~22% better Spark price-performance. They chose transient EMR over EMR Serverless for the predictable nightly job because Spot on a right-sized fleet was ~65% cheaper than Serverless for that steady, large workload — but used Serverless for the unpredictable analyst jobs where no-idle beat Spot pricing. They turned AQE skew-join on after the corporate-account skew turned a 40-minute job into a 5-hour one in week two. They set Athena workgroup scan caps at 3 TB after an analyst’s SELECT * scanned 50 TB.

The outcome (12 months).

Going deeper

The overview is enough to build the platform. This section is for the engineer who has to operate it — the internals, the newer AWS direction, and the edges where things break.

How Lake Formation enforces security inside Spark

The claim that “the same column mask hides msisdn in Athena and in a Spark SELECT on EMR” sounds like magic; the mechanism is credential vending. Normally an engine reads S3 with its own IAM role. Under Lake Formation you register the S3 locations and remove the IAMAllowedPrincipals grant — the backward-compatibility default that lets plain IAM keep working. Leave it in place and your fine-grained policy is decorative, because IAM still grants broad access. Now, when an engine wants data, it does not read S3 directly; it calls Lake Formation (GetDataAccess, and for table access GetTemporaryGlueTableCredentials), Lake Formation checks the caller’s grants, applies column filters and row filters, and hands back short-lived, scoped credentials plus a filtered view of the data. For EMR this rides on the cluster’s runtime role (not its instance profile) and record-server machinery that applies the filtering before rows reach the Spark executors. The upshot: policy lives once in Lake Formation and is enforced at the storage boundary, so it cannot be bypassed by “just reading the bucket.” The wiring rule that makes it real: no human or job role gets direct s3:GetObject on the lake buckets — every read is brokered.

Spot interruption mechanics, and why the node split matters

EC2 Spot capacity is reclaimed when AWS needs it back; you get a two-minute interruption notice (and a rebalance recommendation earlier). On EMR, losing a task node is a non-event — its in-flight Spark tasks are simply rescheduled elsewhere, because task nodes hold no durable state. Losing a core node is different: core nodes carry HDFS blocks and shuffle data, so their loss can force recomputation or, if you lose enough, fail the job. That asymmetry is the entire reason for the core-on-On-Demand, task-on-Spot rule. The price-capacity-optimized allocation strategy (the current best default) picks Spot pools that are both cheap and deep, minimising interruptions versus the older lowest-price strategy, which piles into one cheap-but-shallow pool and gets mass-reclaimed. Diversify the task fleet across many instance types and AZs so no single pool’s reclamation can stall the job. The EC2 Spot with mixed instances and capacity-optimized allocation lesson covers the interruption-handling patterns in depth.

The warehouse edge: Redshift, in depth

“Redshift / Spectrum” in the serving tier hides several distinct choices. RA3 nodes decouple compute from storage: data lives in Redshift Managed Storage (RMS), an S3-backed tier, so you size compute and storage independently and stop over-provisioning nodes just to hold bytes. Redshift Spectrum queries S3 data directly as external tables (via the Glue catalog) without loading it — the warehouse reaches into the lake. Redshift Serverless removes cluster management entirely, billing by RPU-seconds and scaling on demand — the natural fit for spiky BI. And Zero-ETL integrations replicate data from Aurora (MySQL/PostgreSQL), RDS for MySQL, and DynamoDB into Redshift in near-real-time without you building a pipeline — the source keeps serving transactions while an analytics-ready copy stays fresh in the warehouse. Zero-ETL is AWS quietly deleting an entire category of glue code; where it fits, prefer it over hand-built CDC into Redshift.

The newer direction: SageMaker Lakehouse and Unified Studio

The strategic shift worth knowing: AWS is unifying the lake and the warehouse. Amazon SageMaker Lakehouse presents S3 data-lake tables and Redshift warehouse tables through one Iceberg-compatible access layer, governed by Lake Formation — so a single engine can query across lake and warehouse without copying data between them. Amazon SageMaker Unified Studio (announced at re:Invent 2024, generally available in 2025) then wraps data engineering, SQL analytics, ML, and generative AI into one studio over that lakehouse, with Amazon DataZone’s business catalog and governance folded in. You do not have to adopt it to run the architecture in this lesson — but it is the direction the individual services (Glue, Athena, EMR, Redshift, Lake Formation, DataZone) are converging toward, and green-field platforms should expect “lake vs. warehouse” to become one governed surface rather than two.

Search, NL, and the rest of the serving edge

Beyond SQL and BI, two serving patterns recur. Amazon OpenSearch Service is the target when the access pattern is full-text search or log/observability analytics — curated data (or a Firehose branch) flows into OpenSearch for sub-second search and dashboards, complementing rather than replacing Athena’s scan-based SQL; OpenSearch also has a zero-ETL integration to query S3 directly. Amazon Q brings natural language over the top: Q in QuickSight answers “why did revenue dip in APAC” the way an analyst would, and generative SQL in Athena/Redshift turns plain-English questions into queries. These are consumption conveniences layered on the same governed curated tables — they change who can ask, not where the data lives.

Orchestration: three tools, one job

Step Functions, MWAA, and Glue workflows overlap, and picking between them is about blast radius. Glue workflows chain Glue jobs, crawlers, and triggers into a DAG inside Glue — simplest when the whole pipeline is Glue. Step Functions is the AWS-native state machine for event-driven flows that span services (EMR, Glue, Lambda, ECS) with visual retries and error handling — see Step Functions distributed orchestration. MWAA (managed Airflow) wins when you have hundreds of DAGs, complex cross-system dependencies, backfills, and a team already fluent in Airflow’s Python-DAG idiom. Small pipeline: Glue workflows. AWS-native event flow: Step Functions. Big heterogeneous DAG estate: MWAA.

Batch vs. stream: the lambda/kappa framing and where Flink sits

This whole lesson is the batch half — Spark transforming bounded datasets on disposable compute. The complement is stream processing: unbounded event data handled record-by-record at low latency, whose AWS engine is Amazon Managed Service for Apache Flink (formerly Kinesis Data Analytics for Apache Flink). Two classic blueprints frame the choice. Lambda architecture runs both a batch layer (accurate, high-latency — this platform) and a speed layer (fast, approximate — Flink/streaming) and merges them at serving time; it is powerful, but you maintain two codebases. Kappa architecture says “treat everything as a stream,” running a single streaming pipeline and replaying history through the same code when you need a recompute — simpler, but not every heavy transform fits a streaming model. Most enterprises land in between: batch for the heavy, correctness-critical transforms here; streaming on the real-time architecture for sub-second reactions; one shared S3 lake and catalog underneath both.

Quotas, limits, and the edges that bite

A few real limits worth internalising before they surprise you in production. EMR and EMR Serverless consume EC2 vCPU service quotas — a large transient cluster can hit an account vCPU ceiling and simply fail to scale, so raise quotas before the month-end peak, not during it. Athena enforces per-query and per-workgroup limits plus DML concurrency caps; a fleet of analysts can queue behind them, which is a signal to split workgroups. Glue has a default concurrent-run limit per job and DPU quotas. Lake Formation caps the number of filters and tags. And the subtle DR edge from the reliability section bears repeating: the catalog and Lake Formation grants do not live in S3 — they live in Glue/Lake Formation and must be reproduced by IaC in a DR region, or you fail over to a lake of opaque files with no schema and no policy.

When to use it

Use this big-data processing platform when:

Trade-offs and anti-patterns:

Alternatives and how to choose:

Situation Better fit
Sub-second event reaction is the core need (fraud-in-flight, live dashboards) Kinesis + Flink / Lambda — the real-time streaming architecture
The need is storage design — open tables, store-once-query-many governance Data Lakehouse on AWS (Iceberg + Lake Formation + multi-engine)
Pure SQL BI on modest, predictable volume; no Spark/ML Redshift-only warehouse — fewer moving parts
Mostly ad-hoc SQL on S3, no heavy distributed transforms yet Athena + Glue, defer EMR until jobs need Spark muscle
Want a managed lakehouse/Spark platform with less AWS plumbing, multi-cloud Databricks on AWS (Unity Catalog + Delta + managed Spark)
Org runs everything on Kubernetes and wants Spark to share it EMR on EKS as the processing tier of this same architecture

The big-data processing platform on AWS is the right default when heavy, varied batch transformation over a growing dataset would otherwise force you to buy and run peak-sized compute around the clock. By keeping storage permanent on S3 and making compute disposable — transient EMR with Spot for the heavy lifting, EMR Serverless and Glue for the spiky and the routine, Athena for exploration, all governed through one Lake Formation plane — you process any volume on compute sized per job and billed by the second, and the cluster that used to cost the same in July as in December simply ceases to exist.

Practice challenges

Work these in order — they escalate from “read the model” to “make the production call.” Each has a solution with a one-line why.

1 (Beginner). In one sentence each, say which is permanent and which is disposable in this architecture, and name the AWS service for each.

<details><summary>Solution</summary>

Permanent = storage = Amazon S3 (the data-lake spine; it holds raw/processed/curated forever). Disposable = compute = EMR / EMR Serverless / Glue (Spark engines launched per job and terminated after). Why: the entire cost and scaling model follows from decoupling one permanent S3 copy from compute you pay for only while a job runs. </details>

2 (Beginner). You need to run an ad-hoc SELECT COUNT(*) ... WHERE dt='2026-09-01' once to answer a question. Which engine, and why not spin up EMR?

<details><summary>Solution</summary>

Use Amazon Athena — serverless, pay only for the TB scanned, nothing to provision or tear down. Why: a one-off query has zero justification for a Spark cluster; Athena has no standing cost, and partition pruning on dt means it scans one day, not the whole table. </details>

3 (Intermediate). A nightly transient EMR cluster keeps failing when Spot capacity is tight, and when it does complete it sometimes restarts from scratch. The team put all nodes (primary, core, task) on Spot to save money. What is the fix, and what principle did they violate?

<details><summary>Solution</summary>

Put primary and core nodes on On-Demand and keep only task nodes on Spot, with a diversified instance-fleet menu and price-capacity-optimized allocation. Why: core/primary nodes hold HDFS/shuffle data and the application master — losing them to a Spot reclaim kills the job; task nodes are pure compute, so losing one just reruns its tasks. </details>

4 (Intermediate). A Spark stage shows 199 of 200 tasks finishing in seconds and one task running for two hours before the job OOMs. Diagnose it, and give the setting that fixes it without adding nodes.

<details><summary>Solution</summary>

It is data skew — one hot key hashed all its rows into a single shuffle partition on one executor. Enable Adaptive Query Execution with skew-join handling: spark.sql.adaptive.enabled=true and spark.sql.adaptive.skewJoin.enabled=true, so Spark splits the oversized partition at runtime. Why: a bigger cluster only gives the one hot task a bigger executor; only splitting the partition (or broadcasting/salting) parallelises it. </details>

5 (Advanced). Design the compute assignment for four workloads on one lake: (a) a predictable nightly 3 TB ETL, (b) unpredictable analyst Spark jobs a few times a week, © a routine incremental raw→processed load every hour, (d) “what’s in this table?” exploration. Assign an engine to each and justify.

<details><summary>Solution</summary>

(a) Transient EMR on EC2 + Spot — steady, large, predictable, so tuned Spot on a right-sized fleet is cheapest. (b) EMR Serverless — spiky and unpredictable, so scale-to-zero with no sizing beats paying for idle. © AWS Glue with job bookmarks — routine incremental ETL, fastest to author, serverless, Flex for cheap non-urgent runs. (d) Athena — zero standing cost, pay per TB scanned. Why: the engine is a per-workload economic choice; forcing one shape on all four leaves either money or simplicity on the table. </details>

6 (Advanced). After adopting Lake Formation, an auditor asks: “prove that analysts cannot read the msisdn column, including from Spark jobs on EMR, not just in the BI tool.” Explain the mechanism that makes this true and the one configuration mistake that would silently make it false.

<details><summary>Solution</summary>

Register the S3 locations with Lake Formation and grant column-level access; every engine — including EMR via its runtime role — obtains short-lived, scoped credentials from Lake Formation per query, and the column filter is applied at the storage boundary before rows reach the Spark executors, so the mask holds in Spark, Glue, and Athena identically. The silent failure: leaving the IAMAllowedPrincipals default grant in place (or granting any role direct s3:GetObject on the bucket), which lets plain IAM bypass Lake Formation and read the raw column. Why: fine-grained governance is only real if Lake Formation is the sole path to the data. </details>

Common beginner mistakes

These are misconceptions, not error messages — the wrong mental model a beginner brings, and the right one to replace it with.

Glossary

AWSArchitectureEnterpriseReference Architecture
Need this built for real?

Vinod is a Senior Cloud Architect (22+ yrs) — available for Azure / AWS / GCP architecture, landing zones, and migrations.

Work with me

Comments