skip to content

Workflow Orchestration & Transformation

Scheduling pipelines and transforming data once it has landed: Airflow, Prefect and Luigi on the orchestration side, dbt inside the warehouse, plus the managed cloud ETL services. Interviewers ask because 'how does this run every night, and what happens when it fails?' is the actual job.

on this pageshow

explore

questions

169 · 8 sections

In Airflow, what do the >> and << operators do between two tasks in a DAG?

level: juniorimportance: must knowfreq 82%
basics
~10 s

In Airflow, a >> b puts b downstream of a: b is not started until a has finished successfully. The arrow declares execution order only. It moves no data between the tasks.

open as a page

In Airflow, what do the task-level retries and retry_delay arguments control?

level: juniorimportance: must knowfreq 78%
basics
~10 s

In Airflow, retries sets how many extra attempts a task instance gets after its first failure, and retry_delay is the wait between attempts. Between attempts the task sits in the up_for_retry state.

open as a page

In Apache Airflow, what is the difference between an operator and a sensor?

level: juniorimportance: must knowfreq 80%
basics
~20 s

In Airflow an operator is a template for one unit of work — run a Bash command, call a Python function, submit a query. A sensor is an operator that only waits, polling until an external condition becomes true.

open as a page

In Airflow, what does setting catchup=False on a DAG do?

level: juniorimportance: must knowfreq 78%
basics
~20 s

With catchup=False, Airflow schedules only the most recent completed data interval when a DAG becomes active, instead of creating one run for every interval back to start_date. Older intervals are skipped and must be backfilled deliberately.

open as a page

In Airflow, what is an XCom and how does one task push a value another task pulls?

level: juniorimportance: must knowfreq 80%
basics
~10 s

An XCom is Airflow's small cross-task message. A task pushes a keyed value into the metadata database, and a downstream task in the same DAG run pulls it back by task id and key.

open as a page

In Prefect, what is a deployment and how does it differ from running a flow locally?

level: juniorimportance: must knowfreq 78%
basics
~20 s

A Prefect deployment is a server-side record of a flow — entrypoint, code location, default parameters, schedules, work pool and infrastructure overrides — so the API can schedule and trigger runs on remote infrastructure instead of you calling the function yourself.

open as a page

In Prefect, what do the @flow and @task decorators do to a Python function?

level: juniorimportance: must knowfreq 80%
basics
~20 s

@flow turns a function into an orchestrated flow run with tracked state, validated parameters and logging. @task marks a unit of work called inside it, giving each call its own tracked run with retries, caching and concurrency options.

open as a page

In Prefect, what does setting cache_key_fn on a task do?

level: juniorimportance: must knowfreq 55%
basics
~20 s

cache_key_fn computes a string key from the run context and the task's inputs. If a completed task run already recorded that key and it has not expired, Prefect skips the function body and returns the stored result in a Cached state.

open as a page

In Prefect, what is a work pool and what does a worker do with it?

level: middleimportance: must knowfreq 74%
basics
~20 s

A Prefect work pool is a typed queue of scheduled flow runs plus the default infrastructure template for running them. A worker process polls one pool, claims runs, and launches each on that infrastructure — a subprocess, container, or Kubernetes job.

open as a page

In Prefect, how does calling a task with .submit() differ from calling it directly?

level: middleimportance: must knowfreq 68%
basics
~20 s

Calling a Prefect task directly runs it inline in the flow and returns its value, so steps are sequential. Calling .submit() hands it to the flow's task runner and returns a PrefectFuture immediately, letting independent task runs overlap; .result() waits for the value.

open as a page

In Luigi, what determines that a task is already complete and can be skipped?

level: middleimportance: must knowfreq 55%
basics
~20 s

Luigi calls the task's complete() method, which by default returns True only when every Target returned by output() reports exists(). Completion is a property of the data on storage, not of any run history, so reruns skip whatever already exists.

open as a page

In Luigi, what do the requires(), output() and run() methods of a Task define?

level: juniorimportance: should knowfreq 65%
basics
~20 s

A Luigi Task declares its upstream dependencies in requires(), the Target it will produce in output(), and the actual work in run(). Luigi walks requires() to build the graph and checks output() to decide whether run() is needed.

open as a page

What does Luigi not provide that pushed teams toward Airflow-style orchestrators?

level: principalimportance: should knowfreq 48%
basics
~20 s

Luigi has no scheduler of its own, no notion of a scheduled run per data interval, a read-only visualiser and few operational batteries. Teams wanting time-based triggering, catchup, run history and a rich operator ecosystem moved to Airflow.

open as a page

In Luigi, why can a crashed task leave an output that makes the next run skip it?

level: seniorimportance: nice to knowfreq 35%
basics
~20 s

Because Luigi decides completion from Target existence, a half-written file at the output path is indistinguishable from a finished one. The fix is atomic writes: produce the data at a temporary path and move it into place only on clean exit.

open as a page

In dbt, when is the Jinja in a model rendered, and where can you read the resulting SQL?

level: juniorimportance: must knowfreq 70%
basics
~10 s

dbt renders a model's Jinja into plain SQL before anything reaches the warehouse. The rendered SELECT is written to target/compiled/, and the same SQL wrapped in create or merge statements is written to target/run/.

open as a page

How do you define a reusable macro in a dbt project and call it from a model?

level: juniorimportance: must knowfreq 66%
basics
~20 s

Put a {% macro name(args) %} ... {% endmacro %} block in a .sql file under the macros/ directory, then call it from any model with {{ name(args) }}. The macro's rendered text is spliced into the model's SQL at compile time.

open as a page

In dbt, what is the difference between the view and table materializations?

level: juniorimportance: must knowfreq 85%
basics
~20 s

In dbt, a view materialization stores only the SQL, so every read recomputes the query; a table materialization runs the query once per dbt run and stores the rows, making reads fast but data only as fresh as the last build.

open as a page

In dbt, what is a model, and what does a model's .sql file actually contain?

level: juniorimportance: must knowfreq 88%
basics
~20 s

A dbt model is one SELECT statement in a .sql file under models/. dbt renders its Jinja, wraps it in the DDL for the chosen materialization, and builds a table or view named after the file.

open as a page

In dbt, why should a model use ref('stg_orders') instead of a hard-coded table name?

level: juniorimportance: must knowfreq 90%
basics
~20 s

Because ref() does three things a literal name cannot: it declares a dependency so dbt builds the upstream model first, it resolves to the right database and schema for the environment you are running in, and it feeds lineage, docs and selection.

open as a page

What does an AWS Glue crawler create in the Data Catalog after it scans an S3 prefix?

level: juniorimportance: must knowfreq 68%
basics
~20 s

A Glue crawler produces metadata, never data. It creates or updates a table inside a Data Catalog database — inferred columns and types, the S3 location, the format's SerDe — plus one partition entry for each detected sub-prefix.

open as a page

Athena returns zero rows for a partition whose files are in S3 — what went wrong in the Glue Data Catalog?

level: middleimportance: must knowfreq 74%
basics
~20 s

Almost always the partition is not registered. Athena reads only partition entries in the Glue Data Catalog, so a new prefix nobody added is invisible. Re-run the crawler, run MSCK REPAIR TABLE, add the partition explicitly, or use partition projection.

open as a page

In AWS Glue, what does a DynamicFrame give you that a Spark DataFrame does not?

level: middleimportance: must knowfreq 78%
basics
~20 s

A DynamicFrame carries a schema per record instead of one schema for the whole dataset, so conflicting types survive as a ChoiceType rather than failing or being coerced. It also brings Glue-only transforms, Data Catalog reads, error capture and job bookmarks.

open as a page

In an AWS Glue job, what does a job bookmark track between runs?

level: middleimportance: must knowfreq 72%
basics
~20 s

A Glue job bookmark is server-side state, kept per job, recording which input a previous run already processed — for S3 which objects were consumed, for JDBC how far the bookmark key columns advanced — so the next run reads only what is new.

open as a page

How does an AWS Glue crawler decide which parts of an S3 key prefix become partition columns?

level: middleimportance: should knowfreq 52%
basics
~20 s

It walks the folder levels below its target prefix. Hive-style folders named key=value become a partition column called key; any other consistent folder level becomes an unnamed column partition_0, partition_1 and so on, in depth order.

open as a page

In Azure Data Factory, when do you need a Mapping Data Flow instead of a Copy activity?

level: juniorimportance: must knowfreq 78%
basics
~20 s

Use a Copy activity when you only move and land data. Use a Mapping Data Flow when you need real transformation - joins, aggregates, pivots, per-row insert/update logic - built visually and executed on a Spark cluster the service provisions.

open as a page

In Azure Data Factory, how do a linked service, a dataset and an activity differ?

level: juniorimportance: must knowfreq 76%
basics
~20 s

In Azure Data Factory a linked service holds the connection and credentials for a store or compute, a dataset names the data inside it using that linked service, and an activity is the pipeline step that acts on the data.

open as a page

Why does an ADF Mapping Data Flow that processes 500 rows still take several minutes?

level: middleimportance: must knowfreq 55%
basics
~20 s

Nearly all of that time is cluster acquisition. A Data Flow activity must obtain a managed Spark cluster before reading a row. Set a time to live on the Azure integration runtime so later activities reuse a warm cluster.

open as a page

In Azure Data Factory, when would you choose a tumbling window trigger over a schedule trigger?

level: middleimportance: must knowfreq 70%
basics
~20 s

Choose the tumbling window trigger when each run owns a time slice: it fires once per contiguous, non-overlapping window, backfills every elapsed window from its start time, exposes the window boundaries to the pipeline, and supports concurrency limits, retries and dependencies.

open as a page

In an ADF Mapping Data Flow, why does Data preview require debug mode to be on?

level: juniorimportance: should knowfreq 48%
basics
~20 s

Data preview really executes the flow up to that transformation against live data. Debug mode starts a dedicated Spark cluster for your design session; without it there is no compute to run the sampled query on.

open as a page

In Apache Beam, what is a PCollection and what properties does it guarantee?

level: juniorimportance: must knowfreq 80%
basics
~20 s

A PCollection is Apache Beam's immutable, distributed dataset that flows between transforms. It is unordered, either bounded or unbounded, carries a coder and a windowing strategy, and is never modified in place — every PTransform produces a new one.

open as a page

In Apache Beam, why does CombinePerKey handle a hot key better than GroupByKey?

level: middleimportance: must knowfreq 68%
basics
~20 s

GroupByKey sends every value for a key across the shuffle and materializes them all on one worker, so a hot key becomes a straggler or an out-of-memory failure. CombinePerKey pre-aggregates with an associative, commutative CombineFn, so only small accumulators are shuffled and merged.

open as a page

In Dataflow, what is the difference between draining and cancelling a streaming job?

level: middleimportance: must knowfreq 70%
basics
~20 s

Draining a Dataflow streaming job stops it from reading new input, advances the watermark to infinity so every open window fires and in-flight data is written to sinks, then finishes. Cancelling stops immediately and discards buffered state and in-flight data.

open as a page

How do you update a running Dataflow streaming job in place without losing its state?

level: seniorimportance: must knowfreq 60%
basics
~20 s

Submit the new pipeline with the --update flag and the running job's name. Dataflow runs a job-graph compatibility check; if it passes, the replacement job takes over the old job's in-flight state and source position. If it fails, the original keeps running.

open as a page

In an Apache Beam DoFn, what runs in setup versus start_bundle versus process?

level: middleimportance: should knowfreq 58%
basics
~20 s

In an Apache Beam DoFn, setup runs once per DoFn instance on the worker for expensive one-time init, start_bundle once before each bundle of elements, process once per element, finish_bundle after the bundle's elements, and teardown best-effort when the instance is discarded.

open as a page

Why do workflow orchestrators model pipelines as directed acyclic graphs and reject cycles?

level: juniorimportance: must knowfreq 72%
basics
~20 s

An acyclic graph always has a valid execution order, so the scheduler can find work that is ready and can tell when the run is finished. A cycle leaves tasks waiting on each other forever, with no defensible starting point and no termination.

open as a page

What is the difference between ETL and ELT in a data pipeline?

level: juniorimportance: must knowfreq 82%
basics
~20 s

ETL transforms data in a separate processing tier before writing it to the target. ELT loads source-shaped data into the target first and transforms it there, using the target's own compute, keeping the raw input queryable.

open as a page

What makes a scheduled batch task idempotent, and why do orchestrators require it?

level: juniorimportance: must knowfreq 75%
basics
~20 s

A task is idempotent when running it again for the same input window leaves the target in the same state as one successful run. Orchestrators retry, replay and re-run tasks constantly, so anything else duplicates data.

open as a page

What does the pipeline cron schedule `*/15 9-17 * * 1-5` actually fire on?

level: juniorimportance: must knowfreq 70%
basics
~20 s

It fires every fifteen minutes, at :00, :15, :30 and :45, during hours 09 through 17 inclusive, Monday to Friday. That is 36 firings per weekday, the last one at 17:45, and none at weekends.

open as a page

In a workflow orchestrator, when one upstream branch of a fan-in task fails, what decides whether it runs?

level: middleimportance: must knowfreq 68%
basics
~20 s

The downstream task's trigger policy — the rule saying which combination of upstream terminal states lets it start. The default everywhere is all upstreams succeeded, so one failed branch blocks the join and the run ends incomplete rather than publishing partial data.

open as a page