TopicsData engineering
Data engineering
Pipelines, sketches, approximate aggregations, object storage, and stream processing with event time, windows, and exactly-once sinks.
Common tags: sketches, pipelines, aggregations, object-storage, stream-processing
- Data engineering
Workflow Orchestration - Scheduled DAGs vs Durable Execution (Airflow, Temporal & Friends)
Cluster · Workflow Orchestration
Concept hub: two families of orchestrators, scheduled batch DAGs (Airflow, Dagster, Prefect, Argo) vs durable execution (Temporal, Cadence, Step Functions, Durable Functions, Restate, Inngest); the shared kernel (durable state, scheduler, queue, workers, retries, heartbeats/leases, at-least-once units so idempotency matters); comparison tables; runnable task-level vs step-journal crash demo and lease/heartbeat kernel; decision chart.
Open study →- data-engineering
- workflow-orchestration
- airflow
- dags
- dagster
- prefect
- argo-workflows
- scheduling
- backfill
- idempotency
- interview
- Data engineering
Temporal Patterns in Production - Sagas, Child Workflows, Continue-As-New, Versioning & Worker Scaling
Cluster · Workflow Orchestration
Sagas with compensations (register compensation first, non-retryable errors, runnable trip saga), child workflows and parent close policy, continue-as-new and history limits (runnable), versioning with patching vs Worker Versioning (runnable patched-marker and unsafe-change demo), human-in-the-loop with signals/updates and timers, idempotent activities, worker scaling and sticky queues, persistence and visibility stores and history shards; shipping-a-change decision chart.
Open study →- data-engineering
- distributed-systems
- workflow-orchestration
- durable-execution
- temporal
- step-functions
- event-history
- replay
- sagas
- versioning
- interview
- Data engineering
Durable Execution & the Temporal Model - Workflows vs Activities, Event History, Replay & Determinism
Cluster · Workflow Orchestration
Workflows vs activities, event history and replay, determinism rules (SDK time/random, no I/O in workflow code) with a runnable replay engine showing a wall-clock non-determinism bug and the fix; timers, signals, queries and updates (runnable race demo); task queues and workers; retry policy defaults and the four activity timeouts; what exactly-once means (and does not) with idempotency keys; history and payload limits; workflow-vs-activity decision chart.
Open study →- data-engineering
- distributed-systems
- workflow-orchestration
- durable-execution
- temporal
- step-functions
- event-history
- replay
- sagas
- versioning
- interview
- Data engineering
DAG Fundamentals & Airflow Architecture - Topological Scheduling, the Scheduler Loop, Executors, Task States & Pools
Cluster · Workflow Orchestration
Why DAGs and what topological order buys (waves, critical path, cycle detection); Airflow 3.x architecture (Dag processor, Dag bundles, scheduler, API server and Task Execution API, triggerer, metadata DB); task-instance states and trigger rules; pools and concurrency limits; executors compared (Local, Celery, Kubernetes, Edge, multiple executors); HA schedulers via row locks; runnable mini scheduler and wave/mapping simulation; executor decision chart.
Open study →- data-engineering
- workflow-orchestration
- airflow
- dags
- dagster
- prefect
- argo-workflows
- scheduling
- backfill
- idempotency
- interview
- Data engineering
Choosing & Operating Workflow Orchestrators - Airflow vs Dagster vs Prefect vs Temporal vs Step Functions vs Argo, Testing & Migration
Cluster · Workflow Orchestration
Airflow vs Dagster vs Prefect vs Argo vs Temporal vs Step Functions vs cron+queue on run model, graph, state model, latency, duration, dynamic-ness, backfills, human waits, ops burden and cost drivers; what goes wrong if you pick otherwise; observability signals; testing (runnable DAG integrity test and workflow replay-test gate); migrations (cron to Airflow, Airflow 2 to 3, Airflow to Dagster, to durable engines); interview Q&A; decision chart.
Open study →- data-engineering
- workflow-orchestration
- airflow
- dags
- dagster
- prefect
- argo-workflows
- scheduling
- backfill
- idempotency
- interview
- Data engineering
Writing Correct Pipelines - Data Intervals, Catchup & Backfill, Idempotent Tasks, Deferrable Sensors & Assets
Cluster · Workflow Orchestration
Logical date and data intervals, catchup and backfill in Airflow 3 (catchup off by default, logical_date None for asset/API runs), idempotent tasks with partition overwrite (runnable sqlite demo: append vs delete+insert), DST and time zones (runnable America/Chicago demo), XCom limits, sensors vs deferrable operators vs assets and AssetWatcher, dynamic task mapping, Deadline Alerts replacing SLAs, failure modes; Dagster/Prefect comparisons; waiting decision chart.
Open study →- data-engineering
- workflow-orchestration
- airflow
- dags
- dagster
- prefect
- argo-workflows
- scheduling
- backfill
- idempotency
- interview
- Data engineering
Windowing — Tumbling, Hopping, Session & Global Windows with Triggers
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
A window turns an infinite stream into finite chunks you can aggregate. The window type is a product decision, not a tuning knob: tumbling windows answer per minute, hopping windows answer over the last five minutes updated every minute, session windows answer per visit, and a global window answers ever, until you say so. Triggers decide when a window emits, and the accumulation mode decides what each emission means downstream. Getting these wrong produces numbers that look plausible and are silently wrong.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
Stateful Streaming — Keyed State, RocksDB, Checkpoints & Rescaling
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
Stateful streaming means the operator remembers things between events: running counts, open windows, join buffers, dedupe sets, feature rows. That state can be terabytes. It must survive crashes consistently with the source offsets, and it must be redistributable when you change parallelism. Flink's reference design is keyed state partitioned into key groups, stored in RocksDB, snapshotted by barrier checkpoints into object storage, and restorable through savepoints. Spark and Kafka Streams solve the same problem with different trade-offs.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
Stream Processing — Event Time, Windows, State & Exactly-Once
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
Stream processing is continuous computation over unbounded data: you keep results fresh as events arrive instead of recomputing everything on a schedule. The hard part is not reading from Kafka fast. It is answering four questions when data arrives late, out of order, and forever: what you compute, where in event time, when you emit, and how refinements relate. A production engine must also keep large keyed state consistent across crashes and rescaling, and get results into sinks exactly once.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
Streaming Joins, Backpressure & Operations — Interval, Temporal, Skew & Backfill
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
Joins and flow control are where streaming jobs fall over in production. A stream-stream join must buffer both sides, so it needs a time bound or its state grows forever. A stream-table join must decide which version of the row applies: the latest value, or the version that was valid at the event's time. Backpressure is how a healthy engine slows a fast source instead of crashing. Skew is why one subtask is at 100 percent while the rest idle. Operating the job means watching lag, planning a backfill, and upgrading state safely.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
End-to-End Exactly-Once — Replayable Sources, 2PC & Idempotent Sinks
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
Exactly-once in streaming never means each record is physically processed once. After a crash, records are replayed. It means the effect of each record on state and on externally visible output happens once. That property is end-to-end only if three links hold: a replayable source with offsets you can rewind, checkpointed state consistent with those offsets, and a sink that is idempotent or transactional and commits only what a completed checkpoint covers. Break any link and you are back to at-least-once or at-most-once.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
Event Time & Watermarks — Allowed Lateness, Side Outputs & Idle Sources
Cluster · Stream Processing — Event Time, Windows, State & Exactly-Once
In event-time processing the engine never knows a window is complete. A phone that was in airplane mode can upload yesterday's clicks right now. A watermark is the engine's declaration that it believes no more events older than W will arrive. It drives on-time window firing, timers, and state cleanup. Too aggressive and you drop or mis-count late data. Too conservative and every result is delayed. Allowed lateness keeps window state around a bit longer for corrections, a side output captures anything later than that, and idleness handling stops one silent partition from freezing the whole job.
Open study →- data-engineering
- stream-processing
- flink
- spark-structured-streaming
- kafka-streams
- watermarks
- windowing
- exactly-once
- interview
- Data engineering
Security — IAM, Encryption, Public Buckets & Threat Model
Cluster · Object Storage
The famous object-storage incidents are public buckets, over-broad IAM, leaked long-lived keys, and confused-deputy access. This lesson is the threat model: roles instead of static keys, block public access, encryption that does not excuse a public policy, and short-lived sharing.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Object Storage — Buckets, Keys, Consistency & Scale
Cluster · Object Storage
Object storage is the default durable store for lakes, backups, media, and model artifacts: a flat bucket plus key over HTTP, not a POSIX mount. This hub maps block versus file versus object, then consistency, multipart, lifecycle, and the IAM threat model. Metadata partition keys and change streams stay on the sharding and CDC hubs.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Multipart Uploads, Parallelism & Throughput
Cluster · Object Storage
Large objects should not ride one fragile HTTP PUT. Multipart splits bytes into parts, uploads them in parallel, and publishes one object at complete. Interviews expect the state machine, part size, checksums, and abort hygiene so incomplete uploads do not bill forever.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Lifecycle, Storage Classes, CDN & Presigned URLs
Cluster · Object Storage
Capacity is cheap and the wrong tier, request pattern, or egress path is not. Pair storage classes with lifecycle transitions and expiry, put a CDN in front of hot reads with a private origin, and hand browsers a short-lived presigned URL instead of a public bucket.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Data Model — Buckets, Objects, Keys, Versioning & Metadata
Cluster · Object Storage
The object data model is a bucket, bytes plus system metadata, a UTF-8 key, optional user metadata, and optional versions. Interviews fail when prefixes are treated as directories, ETags and version ids are ignored, or one hot prefix throttles PUTs and listings.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Consistency — Read-after-write, Listing & Conditional Writes
Cluster · Object Storage
Modern S3 is strongly consistent for new objects, overwrites, deletes, and listings in-region. Interviews still expect the pre-2020 failure mode, what If-Match buys you, and how replication, CDNs, and client caches put staleness back in front of a strong API.
Open study →- data-engineering
- object-storage
- s3
- blob
- multipart
- lifecycle
- interview
- Data engineering
Exactly-Once CDC Pipelines — Idempotent Consumers, Keys & At-Least-Once Reality
Cluster · CDC & Debezium
CDC capture is at-least-once. Exactly-once effects come from a stable key plus an idempotent sink: an LSN guard for projections, an inbox for side effects. Kafka transactions do not make an email or a charge exactly-once.
Open study →- data-engineering
- cdc
- exactly-once-effects
- idempotent-consumers
- cdc-keys
- tombstones
- debezium
- Data engineering
Debezium & Kafka Connect — Snapshots, Offsets, Schema History & Heartbeats
Cluster · CDC & Debezium
Debezium snapshots a consistent read, then streams from a stored position. Connect offsets are the resume token. Schema history decodes DDL. Heartbeats advance a quiet slot so WAL can be released. Domain events still go through the outbox lesson.
Open study →- data-engineering
- cdc
- debezium
- kafka-connect
- schema-history
- cdc-heartbeats
- wal-tailing
- Data engineering
WAL Tailing vs Query-Based CDC — Log vs Poll Tradeoffs
Cluster · CDC & Debezium
Poll CDC reads a watermark column and misses hard deletes. Log CDC reads commit order from WAL or binlog and keeps a replication slot until the consumer confirms. Pick the log when you can operate it.
Open study →- data-engineering
- cdc
- wal-tailing
- query-cdc
- logical-decoding
- replication-slots
- binlog
- Data engineering
Change Data Capture — WAL Tailing, Debezium & Event Pipelines
Cluster · CDC & Debezium
Dual-write splits one business fact across a database commit and a later publish. CDC reads the database change log so the commit is the event. This hub maps log versus poll, Debezium, the existing outbox lesson, exactly-once effects, and failure modes.
Open study →- data-engineering
- cdc
- change-data-capture
- wal-tailing
- debezium
- pipelines
- dual-write
- Data engineering
CDC Failure Modes — Lag, Schema Breaks, Tombstones, Backfills & Replays
Cluster · CDC & Debezium
CDC fails in production as slot disk, breaking DDL, missing tombstones, and a panicked offset rewind. Freshness is an SLO. Schema changes expand then contract. Rebuild a projection beside the live one, then swap.
Open study →- data-engineering
- cdc
- connector-lag
- schema-breaks
- tombstones
- backfills
- cdc-replays
- Data engineering
T-Digest — Centroid Compression, Tail Accuracy & Heuristic Limits
Cluster · Approximate Aggregations
T-Digest stores mergeable centroids with a compression parameter that spends accuracy on the tails. Interviewers want how centroids work, why p99 looks good, and where the heuristic breaks versus KLL’s theorems.
Open study →- data-engineering
- t-digest
- sketches
- approximate-aggregations
- quantiles
- interview
- Data engineering
Sketch Merge Pipelines — Incremental, Hierarchical & Cross-Shard Aggregation
Cluster · Approximate Aggregations
Sketches only pay off when the pipeline merges them correctly: incrementally on a stream, hierarchically across regions, and across shards without double-counting. Interviewers want leaf-to-region-to-global and the retry failure mode.
Open study →- data-engineering
- merge-pipelines
- sketches
- approximate-aggregations
- stream-processing
- interview
- Data engineering
Production Sketch Ops — Serialization, Idempotent Merges, Bias & Monitoring
Cluster · Approximate Aggregations
Sketches fail in production through bytes, retries, and silent bias — not through forgetting the paper. Version the payload, apply once by sketch_id, and watch coverage plus dual-read error.
Open study →- data-engineering
- sketches
- approximate-aggregations
- merge-pipelines
- observability
- interview
- Data engineering
KLL vs T-Digest vs HDRHistogram vs Exact — When to Choose What
Cluster · Approximate Aggregations
Senior interviews are decision matrices, not brand loyalty. Pick the structure that matches error model, merge needs, value domain, and ops cost — KLL, T-Digest, HDR/Prom hist, or exact offline.
Open study →- data-engineering
- kll
- t-digest
- hdrhistogram
- sketches
- approximate-aggregations
- interview
- Data engineering
KLL Quantile Sketches — Error Bounds, k Parameter & Merge Semantics
Cluster · Approximate Aggregations
KLL (Karnin–Lang–Liberty) is a mergeable quantile sketch with a provable rank-error bound. Interviewers want the k parameter, what ε means, and how merges preserve guarantees — not vendor trivia.
Open study →- data-engineering
- kll
- sketches
- approximate-aggregations
- quantiles
- interview
- Data engineering
Approximate Aggregations — Sketches for Quantiles, Cardinality & Merge Pipelines
Cluster · Approximate Aggregations
Exact p99 over billions of events does not fit in memory, and averaging shard percentiles is wrong. Interviewers expect mergeable sketches: fixed-size summaries you update online, serialize, and combine leaves to regions to global.
Open study →- data-engineering
- approximate-aggregations
- sketches
- kll
- t-digest
- merge-pipelines
- interview