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 pageshowhide
explore
- Apache Airflow37 questions
- DAGs and Tasks6 questions
- Operators and Sensors7 questions
- Scheduling and Triggers6 questions
- XComs and Data Passing6 questions
- Executors and Deployment7 questions
- Monitoring and Retries5 questions
- Prefect18 questions
- Flows and Tasks6 questions
- Deployments and Work Pools6 questions
- States, Caching and Results6 questions
- Luigi4 questions
- dbt38 questions
- Models and Refs6 questions
- Materializations6 questions
- Sources and Seeds7 questions
- Tests and Documentation7 questions
- Jinja and Macros6 questions
- Snapshots and SCD6 questions
- AWS Glue12 questions
- Glue Jobs and DynamicFrames6 questions
- Crawlers and the Data Catalog6 questions
- Azure Data Factory12 questions
- Pipelines, Activities and Triggers6 questions
- Mapping Data Flows6 questions
- Google Cloud Dataflow12 questions
- Apache Beam Model and Runners6 questions
- Streaming Execution and Operations6 questions
- Orchestration Concepts36 questions
- DAGs and Dependency Semantics6 questions
- Scheduling and Data Intervals6 questions
- Idempotency, Backfills and Catchup6 questions
- Retries, SLAs and Alerting6 questions
- Lineage and Pipeline Observability6 questions
- ETL vs ELT and Transformation Placement6 questions
questions
169 · 8 sectionsIn Airflow, what do the >> and << operators do between two tasks in a DAG?
basics
~10 sIn 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.
In Airflow, what do the task-level retries and retry_delay arguments control?
basics
~10 sIn 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.
In Apache Airflow, what is the difference between an operator and a sensor?
basics
~20 sIn 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.
In Airflow, what does setting catchup=False on a DAG do?
basics
~20 sWith 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.
In Airflow, what is an XCom and how does one task push a value another task pulls?
basics
~10 sAn 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.
In Prefect, what is a deployment and how does it differ from running a flow locally?
basics
~20 sA 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.
In Prefect, what do the @flow and @task decorators do to a Python function?
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.
In Prefect, what does setting cache_key_fn on a task do?
basics
~20 scache_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.
In Prefect, what is a work pool and what does a worker do with it?
basics
~20 sA 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.
In Prefect, how does calling a task with .submit() differ from calling it directly?
basics
~20 sCalling 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.
In Luigi, what determines that a task is already complete and can be skipped?
basics
~20 sLuigi 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.
In Luigi, what do the requires(), output() and run() methods of a Task define?
basics
~20 sA 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.
What does Luigi not provide that pushed teams toward Airflow-style orchestrators?
basics
~20 sLuigi 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.
In Luigi, why can a crashed task leave an output that makes the next run skip it?
basics
~20 sBecause 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.
In dbt, when is the Jinja in a model rendered, and where can you read the resulting SQL?
basics
~10 sdbt 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/.
How do you define a reusable macro in a dbt project and call it from a model?
basics
~20 sPut 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.
In dbt, what is the difference between the view and table materializations?
basics
~20 sIn 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.
In dbt, what is a model, and what does a model's .sql file actually contain?
basics
~20 sA 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.
In dbt, why should a model use ref('stg_orders') instead of a hard-coded table name?
basics
~20 sBecause 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.
What does an AWS Glue crawler create in the Data Catalog after it scans an S3 prefix?
basics
~20 sA 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.
Athena returns zero rows for a partition whose files are in S3 — what went wrong in the Glue Data Catalog?
basics
~20 sAlmost 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.
In AWS Glue, what does a DynamicFrame give you that a Spark DataFrame does not?
basics
~20 sA 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.
In an AWS Glue job, what does a job bookmark track between runs?
basics
~20 sA 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.
How does an AWS Glue crawler decide which parts of an S3 key prefix become partition columns?
basics
~20 sIt 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.
In Azure Data Factory, when do you need a Mapping Data Flow instead of a Copy activity?
basics
~20 sUse 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.
In Azure Data Factory, how do a linked service, a dataset and an activity differ?
basics
~20 sIn 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.
Why does an ADF Mapping Data Flow that processes 500 rows still take several minutes?
basics
~20 sNearly 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.
In Azure Data Factory, when would you choose a tumbling window trigger over a schedule trigger?
basics
~20 sChoose 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.
In an ADF Mapping Data Flow, why does Data preview require debug mode to be on?
basics
~20 sData 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.
In Apache Beam, what is a PCollection and what properties does it guarantee?
basics
~20 sA 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.
In Apache Beam, why does CombinePerKey handle a hot key better than GroupByKey?
basics
~20 sGroupByKey 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.
In Dataflow, what is the difference between draining and cancelling a streaming job?
basics
~20 sDraining 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.
How do you update a running Dataflow streaming job in place without losing its state?
basics
~20 sSubmit 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.
In an Apache Beam DoFn, what runs in setup versus start_bundle versus process?
basics
~20 sIn 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.
Why do workflow orchestrators model pipelines as directed acyclic graphs and reject cycles?
basics
~20 sAn 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.
What is the difference between ETL and ELT in a data pipeline?
basics
~20 sETL 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.
What makes a scheduled batch task idempotent, and why do orchestrators require it?
basics
~20 sA 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.
What does the pipeline cron schedule `*/15 9-17 * * 1-5` actually fire on?
basics
~20 sIt 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.
In a workflow orchestrator, when one upstream branch of a fan-in task fails, what decides whether it runs?
basics
~20 sThe 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.