Skip to content
2 changes: 1 addition & 1 deletion projects/backend/beginner/06-cli-user-manager/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ By the end, you should be able to:
```text
Storage: users.json -> [ { id, name, email, role, createdAt } ]

tool add --name "Ada" --email ada@x.com [--role user]
tool add --name "Ada" --email ada@example.com [--role user]
tool list [--role admin] [--sort name]
tool update --id u_01 --name "Ada L."
tool delete --id u_01
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ Ao final, você deve ser capaz de:
```text
Armazenamento: users.json -> [ { id, name, email, role, createdAt } ]

tool add --name "Ada" --email ada@x.com [--role user]
tool add --name "Ada" --email ada@example.com [--role user]
tool list [--role admin] [--sort name]
tool update --id u_01 --name "Ada L."
tool delete --id u_01
Expand Down
113 changes: 83 additions & 30 deletions projects/data-engineering/advanced/01-real-time-platform/README.md
Original file line number Diff line number Diff line change
@@ -1,34 +1,87 @@
# Real-time Data Platform

## Idea
Build a complete real-time data platform supporting streaming ingestion and analytics. Learn about modern data stack architecture.
> 🌐 **English** · [Português](./README.pt-BR.md)

**Domain:** Data Engineering · **Level:** Advanced · **Estimated time:** 1–2 weeks

## Overview

Build an end-to-end real-time data platform: raw events flow in through a durable log, a stream processor turns them into rolling aggregates, and a low-latency serving layer answers dashboard and API queries within seconds of an event landing. Think of the "live metrics" view behind a payments dashboard or a ride-hailing ops screen — the value is entirely in freshness, so a pipeline that is correct but ten minutes stale has failed. This project forces you to reason about the whole path at once: ingestion durability, processing semantics, state, and read-side latency. You will pick where to accept approximation (windowed counts) and where you cannot (money totals), and you will design for the failure modes that only appear when data never stops arriving.

## Prerequisites

- Comfort with a stream processor's core model (Flink, Spark Structured Streaming, or Kafka Streams)
- Experience running a partitioned log like Kafka or Pulsar, including consumer groups and offsets
- A grounding in batch pipelines ([distributed ETL](../02-distributed-etl/) is a useful warm-up)
- Familiarity with windowing, watermarks, and event-time vs processing-time

## Learning Objectives
- Implement streaming infrastructure
- Real-time aggregations
- Build serving layer
- Handle scalability
- Implement monitoring

## Implementation Tips
- Create streaming infrastructure (Kafka, Pulsar)
- Implement stream processing (Flink, Spark)
- Create real-time aggregations
- Build serving layer (Redis, Elasticsearch)
- Implement real-time dashboards
- Add API layer
- Create data quality monitoring
- Implement alerting
- Add latency monitoring
- Create capacity planning
- Implement auto-scaling
- Add disaster recovery
- Build compliance layer
- Create cost optimization

## Key Challenges
- End-to-end latency
- Consistency at scale
- State management
- Infrastructure complexity
- Cost management

By the end, you should be able to:

- Design an ingestion → processing → serving topology with explicit delivery guarantees at each hop
- Choose event-time windowing and watermark strategy for out-of-order and late data
- Manage large keyed state and reason about checkpoint/restore cost
- Separate a hot serving store from the processing layer and justify the split
- Define and measure end-to-end latency and freshness SLOs

## Functional Requirements

1. The platform must ingest events into a partitioned, replayable log and survive a broker restart without data loss.
2. A stream job must compute time-windowed aggregations keyed by a business dimension (e.g. per-merchant, per-region).
3. Late events arriving within a bounded allowed-lateness must still update their window; events beyond it must be routed to a side output, not silently dropped.
4. Results must be written to a serving store that answers point and range queries in single-digit milliseconds.
5. The system must expose end-to-end latency (event timestamp → queryable) as a metric.
6. On job restart from checkpoint, aggregates must not double-count already-processed events.

## Suggested Milestones

1. **Milestone 1 — Ingest & replay:** Stand up the log, produce synthetic events with embedded event-time, and prove you can replay from an offset.
2. **Milestone 2 — Windowed processing:** Implement keyed event-time windows with watermarks, checkpointing, and a late-data side output.
3. **Milestone 3 — Serve & observe:** Sink aggregates to the hot store, add a query API, and instrument freshness and latency SLOs.

## Data & Interface Sketch

```text
producers ─▶ [log: topic "events", N partitions, RF=3]
│ event: {id, merchantId, amountCents, eventTime}
▼
[stream job] keyBy(merchantId)
tumbling 1-min windows, watermark = maxEventTime - 30s
allowedLateness = 5min ─▶ side output "late"
checkpoint every 30s ─▶ durable state backend
│ agg: {merchantId, windowStart, count, sumCents}
▼
[hot store: Redis / key-value] key = merchantId:windowStart
▼
GET /metrics/{merchantId}?from=..&to=.. -> [{windowStart, count, sumCents}]
GET /health/freshness -> { lagSeconds }
```

## Stretch Goals

- Add a second, slower "correction" path (batch reprocessing) and reconcile it against the streaming result — a lambda/kappa comparison.
- Support exactly-once end-to-end by using a transactional sink and idempotent keys.
- Add auto-scaling of processing parallelism driven by consumer lag.

## Definition of Done

- [ ] Events survive a broker or job restart with no loss and no double-counting.
- [ ] Late-but-within-bound events update their window; beyond-bound events land in the side output.
- [ ] The serving API returns windowed aggregates for a key within the target latency.
- [ ] End-to-end freshness lag is exported as a metric and stays under the stated SLO under load.
- [ ] A documented benchmark records throughput, p99 latency, and lag at your target event rate.

## Common Pitfalls

- Mixing processing-time and event-time semantics, so results shift depending on when the job runs.
- Setting watermarks too aggressively and dropping legitimately late data, or too loosely and never closing windows.
- Ignoring state size until checkpoints time out — unbounded keys quietly grow forever.
- Treating the processing store as the serving store, coupling read latency to job restarts.

## Resources

- [Apache Flink: Event Time & Watermarks](https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/) — the canonical model for time in streams.
- [Kafka Documentation: Design](https://kafka.apache.org/documentation/#design) — how the log gives you durability and replay.
- [The Dataflow Model (paper)](https://research.google/pubs/pub43864/) — windowing, watermarks, and triggers, from first principles.
- [Spark Structured Streaming Programming Guide](https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html) — an alternative processing model to compare against.
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
# Plataforma de Dados em Tempo Real

> 🌐 [English](./README.md) · **Português**

**Domínio:** Data Engineering · **Nível:** Avançado · **Tempo estimado:** 1–2 semanas

## Visão Geral

Construa uma plataforma de dados em tempo real de ponta a ponta: eventos brutos chegam por um log durável, um processador de stream os transforma em agregações contínuas, e uma camada de serving de baixa latência responde consultas de dashboards e APIs segundos após o evento chegar. Pense na visão de "métricas ao vivo" por trás de um dashboard de pagamentos ou de uma tela de operações de mobilidade — o valor está inteiramente na atualidade, então um pipeline correto mas dez minutos atrasado falhou. Este projeto força você a raciocinar sobre todo o caminho de uma vez: durabilidade da ingestão, semântica de processamento, estado e latência de leitura. Você escolherá onde aceitar aproximação (contagens em janela) e onde não pode (totais de dinheiro), e projetará para os modos de falha que só aparecem quando os dados nunca param de chegar.

## Pré-requisitos

- Conforto com o modelo central de um processador de stream (Flink, Spark Structured Streaming ou Kafka Streams)
- Experiência operando um log particionado como Kafka ou Pulsar, incluindo grupos de consumidores e offsets
- Uma base em pipelines batch ([ETL distribuído](../02-distributed-etl/) é um bom aquecimento)
- Familiaridade com janelamento, watermarks e event-time vs processing-time

## Objetivos de Aprendizado

Ao final, você deve ser capaz de:

- Projetar uma topologia ingestão → processamento → serving com garantias de entrega explícitas em cada salto
- Escolher janelamento por event-time e estratégia de watermark para dados fora de ordem e atrasados
- Gerenciar estado grande por chave e raciocinar sobre o custo de checkpoint/restore
- Separar um store de serving quente da camada de processamento e justificar a divisão
- Definir e medir SLOs de latência e atualidade de ponta a ponta

## Requisitos Funcionais

1. A plataforma deve ingerir eventos em um log particionado e reproduzível e sobreviver ao reinício de um broker sem perda de dados.
2. Um job de stream deve calcular agregações por janela de tempo, chaveadas por uma dimensão de negócio (ex.: por comerciante, por região).
3. Eventos atrasados que chegarem dentro de uma tolerância de atraso limitada ainda devem atualizar sua janela; eventos além dela devem ser roteados para uma saída lateral, não descartados silenciosamente.
4. Os resultados devem ser gravados em um store de serving que responda consultas pontuais e por intervalo em milissegundos de um dígito.
5. O sistema deve expor a latência de ponta a ponta (timestamp do evento → consultável) como métrica.
6. Ao reiniciar o job a partir de um checkpoint, as agregações não devem contar em dobro eventos já processados.

## Marcos Sugeridos

1. **Marco 1 — Ingerir e reproduzir:** Suba o log, produza eventos sintéticos com event-time embutido e prove que consegue reproduzir a partir de um offset.
2. **Marco 2 — Processamento em janela:** Implemente janelas por event-time chaveadas, com watermarks, checkpointing e uma saída lateral para dados atrasados.
3. **Marco 3 — Servir e observar:** Envie as agregações ao store quente, adicione uma API de consulta e instrumente SLOs de atualidade e latência.

## Esboço de Dados e Interface

```text
produtores ─▶ [log: tópico "events", N partições, RF=3]
│ evento: {id, merchantId, amountCents, eventTime}
▼
[job de stream] keyBy(merchantId)
janelas tumbling de 1min, watermark = maxEventTime - 30s
allowedLateness = 5min ─▶ saída lateral "late"
checkpoint a cada 30s ─▶ backend de estado durável
│ agg: {merchantId, windowStart, count, sumCents}
▼
[store quente: Redis / chave-valor] chave = merchantId:windowStart
▼
GET /metrics/{merchantId}?from=..&to=.. -> [{windowStart, count, sumCents}]
GET /health/freshness -> { lagSeconds }
```

## Desafios Extras

- Adicione um segundo caminho de "correção" mais lento (reprocessamento batch) e reconcilie-o com o resultado do streaming — uma comparação lambda/kappa.
- Suporte exactly-once de ponta a ponta usando um sink transacional e chaves idempotentes.
- Adicione auto-scaling do paralelismo de processamento guiado pelo lag do consumidor.

## Definição de Pronto

- [ ] Eventos sobrevivem ao reinício de um broker ou job sem perda e sem contagem dupla.
- [ ] Eventos atrasados dentro do limite atualizam sua janela; além do limite caem na saída lateral.
- [ ] A API de serving retorna agregações em janela para uma chave dentro da latência alvo.
- [ ] O lag de atualidade de ponta a ponta é exportado como métrica e permanece abaixo do SLO declarado sob carga.
- [ ] Um benchmark documentado registra throughput, latência p99 e lag na sua taxa de eventos alvo.

## Armadilhas Comuns

- Misturar semânticas de processing-time e event-time, fazendo os resultados mudarem conforme o momento em que o job roda.
- Definir watermarks agressivos demais e descartar dados legitimamente atrasados, ou frouxos demais e nunca fechar janelas.
- Ignorar o tamanho do estado até que os checkpoints estourem o tempo — chaves ilimitadas crescem para sempre em silêncio.
- Tratar o store de processamento como store de serving, acoplando a latência de leitura aos reinícios do job.

## Recursos

- [Apache Flink: Event Time e Watermarks](https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/) — o modelo canônico de tempo em streams.
- [Documentação do Kafka: Design](https://kafka.apache.org/documentation/#design) — como o log te dá durabilidade e replay.
- [The Dataflow Model (artigo)](https://research.google/pubs/pub43864/) — janelamento, watermarks e triggers a partir dos princípios.
- [Guia de Programação do Spark Structured Streaming](https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html) — um modelo de processamento alternativo para comparar.
117 changes: 87 additions & 30 deletions projects/data-engineering/advanced/02-distributed-etl/README.md
Original file line number Diff line number Diff line change
@@ -1,34 +1,91 @@
# Distributed ETL System

## Idea
Design a distributed ETL system that scales across multiple nodes. Learn about distributed computing and fault tolerance.
> 🌐 **English** · [Português](./README.pt-BR.md)

**Domain:** Data Engineering · **Level:** Advanced · **Estimated time:** 1–2 weeks

## Overview

Design a distributed ETL job that reads a large dataset, transforms it across many worker nodes, and writes a partitioned result — while staying correct and cheap when a worker dies mid-run. The interesting problems here are not the transformations themselves but everything around them: how data is partitioned, why one hot key can stall a whole stage (data skew), how a shuffle moves gigabytes across the network, and how checkpointing lets you recover without redoing everything. You will treat the cluster as an unreliable machine — nodes vanish, disks fill, one partition is 100× the others — and build a job whose runtime and cost stay predictable anyway. The deliverable is a design and a working distributed job on a framework like Spark, plus a documented plan for skew, retries, and recovery.

## Prerequisites

- Working knowledge of a distributed framework (Apache Spark or Hadoop MapReduce)
- Understanding of partitioning, shuffles, and the map/reduce mental model
- Comfort reasoning about network and disk I/O as the dominant cost
- Familiarity with a columnar format (Parquet/ORC) and object storage

## Learning Objectives
- Implement distributed processing
- Handle distributed state
- Implement fault tolerance
- Optimize resource usage
- Monitor distributed system

## Implementation Tips
- Use distributed framework (Spark, Hadoop)
- Implement distributed tasks
- Handle data distribution
- Create shuffling strategy
- Implement fault tolerance
- Add checkpointing
- Create recovery mechanisms
- Implement monitoring
- Add performance optimization
- Handle skew
- Create resource management
- Implement auto-scaling
- Build operational dashboards
- Create debugging tools

## Key Challenges
- Network I/O optimization
- Data skew handling
- Fault recovery complexity
- Cost optimization
- Debugging distributed issues

By the end, you should be able to:

- Partition input and control parallelism to match cluster resources
- Diagnose and mitigate data skew (salting, broadcast joins, repartitioning)
- Explain what a shuffle does and how to minimize its cost
- Use checkpointing and idempotent writes so a failed run resumes safely
- Instrument a distributed job and read its stage/task metrics to find bottlenecks

## Functional Requirements

1. The job must process a dataset far larger than any single node's memory, in bounded parallelism.
2. Transformations must be deterministic and idempotent so a retried task produces identical output.
3. The job must detect and mitigate skew on at least one join or group-by key.
4. On worker failure, only the lost tasks must re-execute — not the entire job.
5. Output must be written partitioned (e.g. by date) to object storage, with atomic commit so partial writes are never read as complete.
6. The job must emit metrics for records in/out, shuffle bytes, and per-stage duration.

## Suggested Milestones

1. **Milestone 1 — Baseline job:** Read → transform → write partitioned output on a small cluster; confirm correctness on a known sample.
2. **Milestone 2 — Scale & skew:** Run against a large, skewed dataset; measure the straggler, then apply salting or a broadcast join and re-measure.
3. **Milestone 3 — Resilience:** Add checkpointing and atomic/idempotent writes; kill a worker mid-run and verify clean recovery.

## Data & Interface Sketch

```text
[source: object store] raw/events/*.parquet (billions of rows)
│ read, partitions = f(input size, cores)
▼
[map stage] parse, filter, derive columns (no shuffle)
│
▼
[shuffle] repartition by joinKey ── skew? salt hot keys: key -> key#rand(0..N)
│
▼
[reduce stage] join / aggregate (checkpoint here)
│
▼
[sink] write partitioned by dt=YYYY-MM-DD, atomic commit (_SUCCESS marker)
out/agg/dt=2026-07-24/part-*.parquet

Non-functional targets: runtime < T for dataset size S, cost < $C,
recovery re-runs only failed tasks.
```

## Stretch Goals

- Add adaptive query execution (or manual equivalent) that re-partitions based on observed shuffle sizes.
- Support incremental runs that process only new partitions instead of the full dataset.
- Add a spot/preemptible instance pool and prove the job still completes when nodes are reclaimed.

## Definition of Done

- [ ] The job completes on a dataset larger than any node's RAM without OOM.
- [ ] A documented skew mitigation measurably shrinks the slowest task.
- [ ] Killing a worker mid-run reruns only lost tasks and yields identical output.
- [ ] Output is partitioned and only visible after atomic commit; a crashed write leaves no readable partial data.
- [ ] A benchmark records runtime, shuffle bytes, and cost at your target dataset size.

## Common Pitfalls

- Letting the framework pick a default partition count that is far too low or too high for your data.
- Fixing skew by adding memory instead of rebalancing keys — it postpones the failure, it doesn't remove it.
- Non-idempotent writes, so a retried task appends duplicates.
- Reading output before the atomic commit marker exists and treating a partial write as complete.

## Resources

- [Spark: Tuning & Performance](https://spark.apache.org/docs/latest/tuning.html) — memory, serialization, and partitioning guidance.
- [Spark SQL Performance Tuning](https://spark.apache.org/docs/latest/sql-performance-tuning.html) — broadcast joins and adaptive execution for skew.
- [MapReduce (paper)](https://research.google/pubs/pub62/) — the original model for fault-tolerant distributed processing.
- [Apache Parquet documentation](https://parquet.apache.org/docs/) — the columnar format and its partitioning story.
Loading
Loading