High-Performing Data Ingestion and Transformation (Task 3.5)
High-Performing Data Ingestion and Transformation Solutions
Source: https://docs.aws.amazon.com/wellarchitected/latest/performance-efficiency-pillar/Ingestion design is governed by one question — HOW SOON must data be usable? — followed by how MUCH must move. Latency requirements pick the pipeline style; volume and control requirements pick the processing engine.
Ingestion Patterns — Frequency Decides the Tool
| Pattern | Latency | Tool |
|---|---|---|
| Continuous streaming | Seconds | Kinesis Data Streams / MSK |
| Micro-batch delivery | ~60 s | Kinesis Data Firehose |
| Scheduled file transfer | Minutes–hours | DataSync, Transfer Family (SFTP) |
| Database CDC replication | Seconds | DMS |
| Batch upload events | On demand | S3 (direct multipart) |
How to reason it: if action happens on the data in SECONDS (fraud alert, anomaly shutdown), the pipeline must be streaming (Kinesis/MSK) with always-on consumers. If data is merely COLLECTED for later analysis (clickstream to a lake), micro-batch delivery (Firehose) is dramatically simpler and near-enough. Batch windows (nightly ETL) buy the cheapest tooling per byte.
Real use-case: An ad platform first put every impression through a Lambda-per-event pipeline — $18k/month for 30-second freshness nobody consumed. The same stream into Firehose → S3 (Parquet) with hourly Glue aggregation cost $2.4k/month. Freshness was re-specified by the actual consumer (a daily-billed dashboard), and the architecture followed the requirement instead of preceding it.
Gotchas & interview notes: exam cue "sub-second / real-time reactions to events" → Kinesis Data Streams with consumers; "land streaming data in S3/Redshift without managing infrastructure" → Firehose; "partner pushes files nightly over SFTP" → Transfer Family.
The Kinesis Family
- Kinesis Data Streams (KDS): ordered, replayable, partitioned by shard (1 MB/s in, 2 MB/s out per shard); consumers: Lambda, KCL apps, enhanced fan-out (dedicated 2 MB/s per consumer)
- Kinesis Data Firehose: fully managed delivery — no code, auto-scales, transforms via Lambda, buffers to S3 / Redshift / OpenSearch / HTTP; near-real-time loads
- Kinesis Data Analytics: SQL/Flink on the stream (windowed aggregates, alerts)
- Amazon MSK: managed Kafka when you already run Kafka (Kafka APIs/streams)
- S3 zones: raw → cleansed → curated
- AWS Lake Formation: central catalog + fine-grained access control (row/column-level) across the lake — the security answer for data lakes
- Partitioning + Parquet + compaction (small-files problem hurts Athena)
- Redshift Spectrum queries the lake from the warehouse without loading
- Amazon QuickSight (appears as Amazon Quick in the exam guide): serverless BI — SPICE in-memory engine, ML insights/anomaly detection, per-session/user pricing; embeds in apps
- OpenSearch for log search/dashboards; EMR Notebooks/Managed Grafana for ops
- DataSync: agent-based, parallel, checksummed, scheduled — NFS/SMB/HDFS → AWS (up to 10 Gbps/agent)
- Transfer Family: managed SFTP/FTPS/FTP for partners → S3/EFS
- DMS for database CDC ingestion
- 50,000 vehicles stream telemetry (100 events/s each):
- Kinesis Data Streams (100 shards, partition key = vehicleId → per-vehicle ordering); Lambda with enhanced fan-out flags critical faults in <2 s
- Same stream → Firehose buffers to S3
raw/(Parquet via Lambda transform, 5-min window) - Nightly Glue jobs clean/aggregate →
curated/partitioned by date, cataloged in the Glue Data Catalog with Lake Formation granting column-level access (PII columns masked for analysts) - Ad-hoc: Athena; heavy model training: EMR on Spot
- Dashboards: QuickSight (SPICE refresh hourly); DBAs analyze with Redshift Spectrum
Real use-case: IoT fleet of 50k vehicles: partition key = vehicleId → per-vehicle ordering. The fault-detection Lambda uses enhanced fan-out (sub-second, doesn't fight other consumers for the shared 2 MB/s), while a Firehose delivery stream reads the SAME data into S3 for batch analytics — one stream, two consumption models, no duplicate ingestion.
Gotchas & interview notes: "replay 24 h of data" → KDS (retention), never Firehose (no replay — it only delivers forward). "order per key" → partition key design. "multiple consumers without slowing each other" → enhanced fan-out. MSK answers "we already run Kafka / need Kafka ecosystem tools" — otherwise KDS is simpler.
Transformation and Processing
| Tool | Nature | Pick When |
|---|---|---|
| AWS Glue | Serverless Spark ETL + Data Catalog | Scheduled transforms, schema-on-read, no cluster to manage |
| Amazon EMR | Managed Hadoop/Spark/Hive/Presto | You control the cluster, heavy jobs, spot-friendly, custom frameworks |
| AWS Lambda | Function-per-event | Light per-record transforms (Glue/Lambda in Firehose) |
| Amazon Athena | Serverless SQL over S3 | Ad-hoc analytics; pair with Glue catalog |
Glue vs EMR (the control trade): Glue is serverless — no cluster, pay-per-DPU-hour, automatic schema handling in the Data Catalog (one central metadata store for Athena/EMR/Redshift). EMR gives you the cluster: instance types, custom Spark tuning, arbitrary frameworks, Spot fleets — more control, more operations. Heavy, tuned, recurring big-data jobs → EMR; standard ETL and cataloging → Glue.
Format optimization (constant on the exam): convert raw JSON/CSV → Apache Parquet (columnar) and partition by date in S3 — Athena scans 10–100x less data → faster AND cheaper. Athena costs $5/TB scanned.
How Parquet + partitioning pay: a query filtering one column over one day reads ONLY that column of ONLY that partition's files — with JSON row format it would deserialize every byte of the whole table. Columnar + partitioning attacks both CPU (decode less) and I/O (scan less), and cost is proportional to bytes scanned — the SAME optimization is a performance AND a cost answer.
Real use-case: A 90-day log lake in JSON: daily Athena query = $42 (scanning 8.4 TB). Glue job converts to Parquet partitioned by date: the same query scans 60 GB = $0.30. One-time transform, 140x recurring savings — the single highest-leverage optimization in lake design.
Gotchas & interview notes: "reduce Athena cost without changing queries" → Parquet + partitioning. Small-files problem: many tiny Parquet files slow Athena (per-file overhead) — compact them. Redshift Spectrum = query the lake from Redshift WITHOUT loading it (separate from Athena by warehouse context).
Data Lake Construction
- Kinesis Data Streams (KDS): ordered, replayable, partitioned by shard (1 MB/s in, 2 MB/s out per shard); consumers: Lambda, KCL apps, enhanced fan-out (dedicated 2 MB/s per consumer)
- Kinesis Data Firehose: fully managed delivery — no code, auto-scales, transforms via Lambda, buffers to S3 / Redshift / OpenSearch / HTTP; near-real-time loads
- Kinesis Data Analytics: SQL/Flink on the stream (windowed aggregates, alerts)
- Amazon MSK: managed Kafka when you already run Kafka (Kafka APIs/streams)
- S3 zones: raw → cleansed → curated
- AWS Lake Formation: central catalog + fine-grained access control (row/column-level) across the lake — the security answer for data lakes
- Partitioning + Parquet + compaction (small-files problem hurts Athena)
- Redshift Spectrum queries the lake from the warehouse without loading
- Amazon QuickSight (appears as Amazon Quick in the exam guide): serverless BI — SPICE in-memory engine, ML insights/anomaly detection, per-session/user pricing; embeds in apps
- OpenSearch for log search/dashboards; EMR Notebooks/Managed Grafana for ops
- DataSync: agent-based, parallel, checksummed, scheduled — NFS/SMB/HDFS → AWS (up to 10 Gbps/agent)
- Transfer Family: managed SFTP/FTPS/FTP for partners → S3/EFS
- DMS for database CDC ingestion
- 50,000 vehicles stream telemetry (100 events/s each):
- Kinesis Data Streams (100 shards, partition key = vehicleId → per-vehicle ordering); Lambda with enhanced fan-out flags critical faults in <2 s
- Same stream → Firehose buffers to S3
raw/(Parquet via Lambda transform, 5-min window) - Nightly Glue jobs clean/aggregate →
curated/partitioned by date, cataloged in the Glue Data Catalog with Lake Formation granting column-level access (PII columns masked for analysts) - Ad-hoc: Athena; heavy model training: EMR on Spot
- Dashboards: QuickSight (SPICE refresh hourly); DBAs analyze with Redshift Spectrum
Visualization
- Kinesis Data Streams (KDS): ordered, replayable, partitioned by shard (1 MB/s in, 2 MB/s out per shard); consumers: Lambda, KCL apps, enhanced fan-out (dedicated 2 MB/s per consumer)
- Kinesis Data Firehose: fully managed delivery — no code, auto-scales, transforms via Lambda, buffers to S3 / Redshift / OpenSearch / HTTP; near-real-time loads
- Kinesis Data Analytics: SQL/Flink on the stream (windowed aggregates, alerts)
- Amazon MSK: managed Kafka when you already run Kafka (Kafka APIs/streams)
- S3 zones: raw → cleansed → curated
- AWS Lake Formation: central catalog + fine-grained access control (row/column-level) across the lake — the security answer for data lakes
- Partitioning + Parquet + compaction (small-files problem hurts Athena)
- Redshift Spectrum queries the lake from the warehouse without loading
- Amazon QuickSight (appears as Amazon Quick in the exam guide): serverless BI — SPICE in-memory engine, ML insights/anomaly detection, per-session/user pricing; embeds in apps
- OpenSearch for log search/dashboards; EMR Notebooks/Managed Grafana for ops
- DataSync: agent-based, parallel, checksummed, scheduled — NFS/SMB/HDFS → AWS (up to 10 Gbps/agent)
- Transfer Family: managed SFTP/FTPS/FTP for partners → S3/EFS
- DMS for database CDC ingestion
- 50,000 vehicles stream telemetry (100 events/s each):
- Kinesis Data Streams (100 shards, partition key = vehicleId → per-vehicle ordering); Lambda with enhanced fan-out flags critical faults in <2 s
- Same stream → Firehose buffers to S3
raw/(Parquet via Lambda transform, 5-min window) - Nightly Glue jobs clean/aggregate →
curated/partitioned by date, cataloged in the Glue Data Catalog with Lake Formation granting column-level access (PII columns masked for analysts) - Ad-hoc: Athena; heavy model training: EMR on Spot
- Dashboards: QuickSight (SPICE refresh hourly); DBAs analyze with Redshift Spectrum
Transfer and Batch Movement
- Kinesis Data Streams (KDS): ordered, replayable, partitioned by shard (1 MB/s in, 2 MB/s out per shard); consumers: Lambda, KCL apps, enhanced fan-out (dedicated 2 MB/s per consumer)
- Kinesis Data Firehose: fully managed delivery — no code, auto-scales, transforms via Lambda, buffers to S3 / Redshift / OpenSearch / HTTP; near-real-time loads
- Kinesis Data Analytics: SQL/Flink on the stream (windowed aggregates, alerts)
- Amazon MSK: managed Kafka when you already run Kafka (Kafka APIs/streams)
- S3 zones: raw → cleansed → curated
- AWS Lake Formation: central catalog + fine-grained access control (row/column-level) across the lake — the security answer for data lakes
- Partitioning + Parquet + compaction (small-files problem hurts Athena)
- Redshift Spectrum queries the lake from the warehouse without loading
- Amazon QuickSight (appears as Amazon Quick in the exam guide): serverless BI — SPICE in-memory engine, ML insights/anomaly detection, per-session/user pricing; embeds in apps
- OpenSearch for log search/dashboards; EMR Notebooks/Managed Grafana for ops
- DataSync: agent-based, parallel, checksummed, scheduled — NFS/SMB/HDFS → AWS (up to 10 Gbps/agent)
- Transfer Family: managed SFTP/FTPS/FTP for partners → S3/EFS
- DMS for database CDC ingestion
- 50,000 vehicles stream telemetry (100 events/s each):
- Kinesis Data Streams (100 shards, partition key = vehicleId → per-vehicle ordering); Lambda with enhanced fan-out flags critical faults in <2 s
- Same stream → Firehose buffers to S3
raw/(Parquet via Lambda transform, 5-min window) - Nightly Glue jobs clean/aggregate →
curated/partitioned by date, cataloged in the Glue Data Catalog with Lake Formation granting column-level access (PII columns masked for analysts) - Ad-hoc: Athena; heavy model training: EMR on Spot
- Dashboards: QuickSight (SPICE refresh hourly); DBAs analyze with Redshift Spectrum
Worked Example: IoT Fleet Analytics at Scale
- Kinesis Data Streams (KDS): ordered, replayable, partitioned by shard (1 MB/s in, 2 MB/s out per shard); consumers: Lambda, KCL apps, enhanced fan-out (dedicated 2 MB/s per consumer)
- Kinesis Data Firehose: fully managed delivery — no code, auto-scales, transforms via Lambda, buffers to S3 / Redshift / OpenSearch / HTTP; near-real-time loads
- Kinesis Data Analytics: SQL/Flink on the stream (windowed aggregates, alerts)
- Amazon MSK: managed Kafka when you already run Kafka (Kafka APIs/streams)
- S3 zones: raw → cleansed → curated
- AWS Lake Formation: central catalog + fine-grained access control (row/column-level) across the lake — the security answer for data lakes
- Partitioning + Parquet + compaction (small-files problem hurts Athena)
- Redshift Spectrum queries the lake from the warehouse without loading
- Amazon QuickSight (appears as Amazon Quick in the exam guide): serverless BI — SPICE in-memory engine, ML insights/anomaly detection, per-session/user pricing; embeds in apps
- OpenSearch for log search/dashboards; EMR Notebooks/Managed Grafana for ops
- DataSync: agent-based, parallel, checksummed, scheduled — NFS/SMB/HDFS → AWS (up to 10 Gbps/agent)
- Transfer Family: managed SFTP/FTPS/FTP for partners → S3/EFS
- DMS for database CDC ingestion
- 50,000 vehicles stream telemetry (100 events/s each):
- Kinesis Data Streams (100 shards, partition key = vehicleId → per-vehicle ordering); Lambda with enhanced fan-out flags critical faults in <2 s
- Same stream → Firehose buffers to S3
raw/(Parquet via Lambda transform, 5-min window) - Nightly Glue jobs clean/aggregate →
curated/partitioned by date, cataloged in the Glue Data Catalog with Lake Formation granting column-level access (PII columns masked for analysts) - Ad-hoc: Athena; heavy model training: EMR on Spot
- Dashboards: QuickSight (SPICE refresh hourly); DBAs analyze with Redshift Spectrum