Courseiva

CCNA Pde Maintaining Automating Questions

75 of 83 questions · Page 1/2 · Pde Maintaining Automating topic · Answers revealed

1
Multi-Selecthard

You are setting up Dataplex data quality rules for a BigQuery table. You want to define rules that check for non-null values in key columns and also validate that a column's values fall within a certain range. Which TWO rule types must you use? (Choose 2)

Select 2 answers
A.Table rule (e.g., row count)
B.Row rule (e.g., not null)
C.Partition rule
D.Column rule (e.g., value range)
E.Custom SQL rule
AnswersB, D

A row rule evaluates each record independently, so a not-null condition on key columns is expressed as a row-level rule. Range validation on a column's values is also row-scoped, checking every value against the defined bounds. Both satisfy the non-null and range checks required.

Why this answer

Dataplex data quality rules include row rules (for null checks) and column rules (for range or value checks). Table rules apply to the entire table; custom SQL can be used but row and column rules are the standard.

2
MCQmedium

Your team runs a Cloud Composer (Airflow) environment that executes a nightly BigQuery ELT DAG. A downstream task must only run after an upstream task that loads a partitioned table completes, and you want the downstream task to wait for the load's completion signal without polling BigQuery repeatedly. Which Airflow mechanism should you implement to coordinate these tasks within the DAG?

A.Define a direct downstream dependency (loader >> downstream) so Airflow's scheduler marks the downstream task runnable only after the loader task succeeds.
B.Use a task dependency with an ExternalTaskSensor pointed at the loader task's DAG and task ID, configured with the correct execution_delta.
C.Set the loader task's trigger_rule to all_done and add a dependency edge so the downstream task runs immediately after the loader reaches any terminal state.
D.Use an XCom to push a completion flag from the loader task and have the downstream task pull it via xcom_pull, combined with a trigger rule of all_success on the dependency edge.
AnswerA

In Airflow, a dependency edge makes the scheduler hold the downstream task until the upstream task reaches a successful terminal state. This is the native, no-polling coordination mechanism inside a single DAG, exactly matching the requirement that the downstream task wait for the load's completion signal without repeatedly querying BigQuery.

Why this answer

A direct dependency edge is the canonical Airflow way to sequence tasks inside one DAG: the scheduler keeps the downstream task in a waiting state until the upstream load succeeds. Sensors and XComs serve other purposes and would add complexity or incorrect semantics here. The dependency edge provides the required completion gating with no polling of BigQuery.

Exam trap

The trap here is assuming a sensor or XCom is needed when a simple dependency edge already enforces ordering within the same DAG.

3
MCQeasy

You need to estimate the cost of a BigQuery query before running it. Which command or feature should you use?

A.Check the BigQuery jobs list for similar queries.
B.Use the BigQuery cache to estimate if the query is cached.
C.Run EXPLAIN on the query to see the query plan.
D.Use the bq command with the --dry_run flag.
AnswerD

The --dry_run flag validates the query and returns the bytes it would process without executing it, letting you estimate on-demand cost from that byte count. This satisfies the stem's requirement to price a query before running it.

Why this answer

The `bq` command with the `--dry_run` flag allows you to estimate the amount of data a BigQuery query will process before actually executing it. This dry run does not read any data or incur charges; it simply returns the estimated bytes to be processed, which you can use to calculate the cost based on BigQuery's pricing model.

Exam trap

A common trap is confusing the EXPLAIN command (which shows the query plan) with the `--dry_run` flag (which estimates bytes processed and cost).

How to eliminate wrong answers

Option A is wrong because checking the BigQuery jobs list for similar queries only gives you historical cost data, not an estimate for the specific query you are about to run, and it assumes a similar query exists. Option B is wrong because the BigQuery cache stores results of previously run queries, but it does not provide an estimate of cost or data processed; it only indicates whether results might be served from cache. Option C is wrong because running EXPLAIN on the query shows the query plan and execution steps, but it does not provide a cost estimate or the amount of data that will be scanned.

4
MCQmedium

Your team runs a Cloud Composer 2 environment in project `analytics-prod`. A nightly DAG loads Cloud Storage files into BigQuery, but a recent Cloud Storage outage caused several tasks to fail after exhausting their retries. You want failed task instances to automatically re-run without manual intervention once the upstream dependency recovers. What should you do?

A.Set the Airflow configuration option `[core] default_task_retries` to a higher value in the environment configuration and restart the schedulers.
B.Increase the DAG's `retry_delay` and `execution_timeout` so tasks wait longer for Cloud Storage to recover before failing.
C.Configure the DAG with `catchup=True` and set `max_active_runs=1` so missed schedule intervals are backfilled.
D.Define a task-level `on_failure_callback` that calls the Airflow REST API to clear the failed task instances after a delay.
AnswerD

Clearing failed task instances through the Airflow REST API re-queues them so the scheduler runs them again once the upstream Cloud Storage dependency is healthy. Wrapping this in an `on_failure_callback` with a delay automates recovery without manual operator intervention, which is exactly the requirement after retries are exhausted.

Why this answer

Failed task instances that have exhausted retries stay in a failed state until they are cleared, at which point the scheduler re-queues them. Triggering that clear programmatically, for example from a failure callback that waits for the dependency to recover, restores the pipeline automatically. Backfill settings, retry counts, and timeout tuning only affect future scheduling or in-flight attempts, not instances already marked failed.

Exam trap

The trap here is assuming that increasing retry counts or delays will recover tasks that have already reached the failed state, when only clearing and re-queuing those task instances makes the scheduler run them again.

5
MCQmedium

Your Dataflow streaming pipeline writes to BigQuery and occasionally fails with `QuotaExceededError` on streaming inserts during peak hours. You want to reduce insert-driven quota pressure without changing the pipeline's output schema or losing exactly-once semantics. Which change should you make?

A.Add a `GroupByKey` before the BigQuery sink to batch records into larger write requests.
B.Switch the pipeline from the BigQueryIO streaming insert method to the Storage Write API with exactly-once semantics.
C.Increase the number of Dataflow workers and set `--maxNumWorkers` higher so inserts are spread across more threads.
D.Enable `--enableStreamingEngine` and raise the pipeline's `--diskSizeGb` to buffer failed inserts locally.
AnswerB

The Storage Write API uses a different quota pool than legacy streaming inserts and supports exactly-once delivery when configured with the appropriate commit strategy. Migrating the sink reduces pressure on streaming insert quotas while preserving the schema and the exactly-once guarantee the pipeline relies on, directly addressing the peak-hour failures.

Why this answer

Legacy BigQuery streaming inserts draw on a quota pool that is easy to saturate at peak throughput. The Storage Write API is a separate ingestion path with higher and separately metered limits, and its exactly-once mode uses write streams with offsets so retries do not duplicate rows. Switching the sink preserves the output schema while relieving the quota bottleneck.

Exam trap

The trap here is assuming that scaling Dataflow workers or batching requests raises the BigQuery insert quota, when the quota is enforced server-side and only a different ingestion API changes which limit applies.

6
MCQmedium

You are designing a data quality pipeline that must inspect PII in BigQuery tables and de-identify sensitive columns before sharing with analysts. Which GCP service should you use?

A.Dataplex
B.Cloud Data Catalog
C.Cloud DLP
D.Dataflow
AnswerC

Cloud DLP natively scans BigQuery tables, using infoType detectors to identify PII and de-identification transforms such as masking, tokenisation, and format-preserving encryption to protect sensitive columns. This directly satisfies the stem's requirement to inspect and de-identify PII in place before analysts access the shared data.

Why this answer

Cloud DLP (Data Loss Prevention) is the correct choice because it is purpose-built for inspecting, classifying, and de-identifying sensitive data such as PII. It integrates natively with BigQuery via inspection jobs and de-identification templates, allowing you to scan tables for over 150 built-in infoTypes (e.g., email, SSN) and apply transformations like masking, tokenization, or encryption before sharing data with analysts.

Exam trap

The trap here is that candidates often confuse Dataplex's data governance features (like policy tags and metadata) with actual de-identification, but Dataplex cannot transform data—it only applies access controls, whereas Cloud DLP performs the actual masking or tokenization of sensitive values.

How to eliminate wrong answers

Option A is wrong because Dataplex is a data fabric service for managing, governing, and cataloging data across lakes and warehouses, but it does not perform de-identification or PII inspection itself; it can integrate with Cloud DLP for such tasks but is not the primary tool. Option B is wrong because Cloud Data Catalog is a metadata management service for discovering and tagging assets, but it lacks native de-identification capabilities and cannot transform sensitive data. Option D is wrong because Dataflow is a stream/batch processing service that can be used to build custom de-identification pipelines, but it requires manual implementation of DLP logic and is not the out-of-the-box service for inspecting and de-identifying PII in BigQuery tables.

7
MCQeasy

You need to schedule a simple workflow that fetches data from an API every hour, transforms it using Cloud Functions, and writes the result to Cloud Storage. The workflow has no complex branching or retry logic beyond basic retries. Which orchestration service is the MOST cost-effective and simplest to implement?

A.Cloud Scheduler
B.Workflows
C.Cloud Composer
D.Dataflow
AnswerB

Workflows is serverless and billed per step executed, so an hourly fetch-transform-write sequence with basic retries costs far less than provisioning Cloud Composer's always-running Airflow environment, and its YAML definition is simpler for linear, non-branching orchestration.

Why this answer

Workflows is a serverless, fully managed orchestration service that charges only per step executed, making it ideal for simple linear workflows with basic retries. It natively integrates with Cloud Functions and Cloud Storage via connectors, so the hourly fetch-transform-write pipeline can be defined in YAML/JSON without managing infrastructure. Cloud Scheduler alone can trigger jobs but cannot orchestrate multi-step logic or handle the transform step.

Exam trap

PDE often tests the distinction between a scheduler (Cloud Scheduler), an orchestrator (Workflows/Composer), and a data processor (Dataflow); candidates wrongly pick Cloud Scheduler because the question mentions 'every hour' and overlook the multi-step orchestration requirement.

How to eliminate wrong answers

Option A is wrong because Cloud Scheduler is only a cron-based trigger service — it can invoke a function hourly but cannot orchestrate the multi-step fetch-transform-write sequence or manage state between steps. Option C is wrong because Cloud Composer is a managed Apache Airflow service designed for complex DAGs with branching, dependencies, and heavy retry logic; it carries significant cost (environment minimums) and operational overhead that is overkill for a simple linear workflow. Option D is wrong because Dataflow is a stream/batch data processing service based on Apache Beam, not an orchestration engine — it processes data but does not schedule or coordinate API calls and Cloud Functions.

8
Multi-Selecthard

A company runs a Dataflow pipeline that processes a high-volume data stream. They notice that the pipeline's worker CPU utilisation is near 100% and the system lag is increasing. Which three actions can improve performance? (Choose three.)

Select 3 answers
A.Increase the worker disk size.
B.Increase the number of workers.
C.Use batch processing instead of streaming.
D.Enable Dataflow Streaming Engine.
E.Use higher-CPU machine types (e.g., n2-highcpu).
AnswersB, D, E

More workers distribute the load and reduce CPU per worker.

Why this answer

Increasing the number of workers distributes the processing load across more parallel workers, reducing CPU utilization per worker and allowing the pipeline to keep up with the incoming data stream. This directly addresses both high CPU usage and increasing system lag by scaling out horizontally.

Exam trap

A common trap is assuming that increasing disk size (Option A) improves CPU-related performance issues. In Dataflow, increasing disk size only helps with storage bottlenecks (e.g., shuffle disk overflow), not with high CPU utilization or system lag.

9
MCQmedium

Your team runs a Cloud Composer 2 environment to orchestrate BigQuery ELT jobs. A DAG that loads a critical fact table must only run after an upstream ingestion DAG completes, and you want Composer itself to trigger it automatically without any external scheduler. Which mechanism should you use?

A.A TimeDeltaSensor with a fixed delay that matches the upstream DAG's typical runtime.
B.A Pub/Sub push subscription that invokes a Cloud Function to call the downstream DAG's trigger endpoint.
C.An ExternalTaskSensor that watches the upstream DAG's task in the same Composer environment.
D.A Cloud Scheduler job that calls the Airflow REST API to trigger the downstream DAG on a cron schedule.
AnswerC

ExternalTaskSensor is designed precisely for cross-DAG dependencies within the same Airflow deployment. It polls the metadata database for the specified external DAG and task run in the matching logical date, then releases downstream tasks only when that upstream task reaches the expected state. This gives a real completion dependency instead of guessing at timing.

Why this answer

Airflow's ExternalTaskSensor exists specifically to express a dependency on a task in another DAG within the same Airflow deployment. It checks the external DAG run for the matching execution date and only allows downstream tasks to proceed after the upstream task succeeds. Time-based waits, external cron triggers, and message-based glue all fail to guarantee that the ingestion actually finished successfully before the load begins.

Exam trap

The trap here is assuming any scheduling mechanism that can trigger a DAG also enforces a true completion dependency, when only a sensor that inspects the upstream task state does.

10
MCQmedium

You are using Cloud Composer to orchestrate a data pipeline that runs a Dataproc job to process data, followed by a BigQuery load. You notice that the Dataproc job sometimes takes longer than expected, causing the BigQuery load to start before the Dataproc job finishes, resulting in incomplete data. Which Airflow feature should you use to ensure the BigQuery load only runs after the Dataproc job completes successfully?

A.Set a dependency between the Dataproc job task and the BigQuery load task using the >> operator.
B.Use a TimeSensor to wait for the Dataproc job to finish.
C.Use the trigger_rule parameter to set the BigQuery load task to 'all_done'.
D.Set the Dataproc job task's retries to a high number.
AnswerA

In Airflow, task dependencies are defined using the bitshift operators, such as >>, to specify that one task must complete successfully before another starts. By setting the Dataproc job task upstream of the BigQuery load task, you ensure that the load only runs after the Dataproc job finishes successfully. This is the fundamental way to control execution order in a DAG and directly addresses the issue of premature BigQuery loads due to timing.

Why this answer

The correct way to enforce that the BigQuery load runs only after the Dataproc job completes successfully is to define a task dependency using the >> operator. This creates a directed edge in the DAG, ensuring the downstream task waits for the upstream task to succeed. Other options either do not control ordering or do not enforce success.

Task dependencies are the core mechanism in Airflow for sequencing tasks in a data pipeline.

Exam trap

The trap here is confusing task dependencies with sensors or retries; only explicit dependencies guarantee that one task waits for another's successful completion.

11
MCQmedium

A retail analytics team runs a Cloud Composer DAG that starts a Dataproc job, waits for it, and then runs a BigQuery load. The Dataproc job sometimes fails due to a transient YARN resource shortage, and the team wants the DAG to automatically retry the Dataproc submission a few times before alerting. Which Airflow configuration on the Dataproc task best meets this need?

A.Set the Dataproc operator's cluster_name parameter to a different cluster and rely on Dataproc to resubmit the job.
B.Add a downstream EmailOperator that notifies the team on failure and manually reruns the DAG.
C.Set retries to 3 and retry_delay to 5 minutes on the Dataproc job task.
D.Set max_active_runs to 3 on the DAG so multiple DAG runs can absorb the failure.
AnswerC

Task-level retries cause Airflow to re-run the same operator up to the specified count after a failure, with the given delay between attempts. Because the failure is transient, a small number of retries with a short delay gives the cluster time to free resources, and only after all attempts fail does the task enter a failed state that triggers alerting.

Why this answer

Transient failures are best handled with operator-level retries, which re-execute the same task after a delay and only surface an alert once the configured attempts are exhausted. Setting retries to 3 with a short retry_delay on the Dataproc task gives the cluster room to recover from resource pressure while keeping the DAG's dependency chain and alerting behavior unchanged.

Exam trap

The trap here is reaching for DAG-level concurrency settings to solve a single-task reliability problem, when retries belong on the failing task itself.

12
MCQeasy

You need to schedule a Dataproc Spark job to run at 2 AM every day, and upon completion, trigger a BigQuery load job. Which Cloud Composer operator should you use to run the Spark job?

A.DataflowPythonOperator
B.BigQueryOperator
C.DataprocClusterCreateOperator
D.DataprocSubmitJobOperator
AnswerD

DataprocSubmitJobOperator submits a Spark job to a Dataproc cluster and waits for completion, so it fits the scheduled 2 AM run. Its downstream task can then trigger the BigQuery load job within the same DAG.

Why this answer

The DataprocSubmitJobOperator is specifically designed to submit a job (e.g., a Spark job) to an existing Dataproc cluster. In this scenario, you need to run a Spark job on a scheduled basis, and Cloud Composer (Airflow) provides this operator to submit the job to Dataproc. After the Spark job completes, you can chain a BigQuery load operator to trigger the load, matching the requirement exactly.

Exam trap

The trap here is that candidates confuse operators that manage cluster lifecycle (like DataprocClusterCreateOperator) with operators that submit jobs, or they mistakenly think DataflowPythonOperator can run Spark jobs because both are data processing frameworks.

How to eliminate wrong answers

Option A is wrong because DataflowPythonOperator is used to run Apache Beam pipelines on Dataflow, not Spark jobs on Dataproc. Option B is wrong because BigQueryOperator is used to execute BigQuery SQL queries or load jobs, not to run Spark jobs. Option C is wrong because DataprocClusterCreateOperator is used to create a new Dataproc cluster, not to submit a job to an existing cluster; the question assumes the cluster already exists or is managed separately, and the focus is on submitting the Spark job.

13
MCQmedium

A data engineering team manages BigQuery datasets across multiple projects. They want to automatically detect and respond when a scheduled query fails, and they want the response to create an incident in their existing ticketing system. The team prefers minimal custom infrastructure and wants to use native Google Cloud tooling. Which approach should they use?

A.Schedule a Cloud Scheduler job that polls the BigQuery INFORMATION_SCHEMA.JOBS view and emails the team when errors are found.
B.Configure BigQuery to send failure notifications to a Cloud Storage bucket and use a Dataproc job to parse them.
C.Enable BigQuery audit logs, create a log-based alerting policy in Cloud Monitoring, and route the alert to a Pub/Sub topic that triggers a Cloud Function to open the ticket.
D.Use the BigQuery REST API from a long-running Compute Engine VM to poll for failed jobs and call the ticketing API.
AnswerC

BigQuery writes job failure information to Cloud Audit Logs. A log-based alerting policy in Cloud Monitoring can match failed query jobs and publish to a Pub/Sub topic, which invokes a Cloud Function that calls the ticketing API. This uses native tooling with minimal custom infrastructure and reliably captures scheduled query failures.

Why this answer

Cloud Audit Logs capture BigQuery job failures, and log-based alerting policies in Cloud Monitoring can trigger Pub/Sub notifications. A Cloud Function subscribed to the topic can call the ticketing API, creating incidents automatically. This event-driven chain uses managed services, avoids polling, and requires little custom code compared with the other options.

Exam trap

The trap here is choosing polling-based approaches, which add latency and maintenance, instead of event-driven log-based alerting that natively captures BigQuery failures.

14
MCQhard

You manage a BigQuery reservation with 500 baseline slots and autoscaling up to 2000 slots. Your team runs a mix of interactive queries and batch load jobs. During peak hours, you notice that interactive queries are throttled when autoscaling slots are consumed by long-running batch loads. How can you ensure interactive queries get priority access to slots?

A.Create a separate reservation for interactive queries with a higher priority assignment.
B.Reduce the baseline slots to 200 and rely solely on autoscaling.
C.Switch to on-demand pricing to eliminate slot contention.
D.Set the autoscaling max to 1000 slots for batch jobs.
AnswerA

Separate reservations isolate slot pools, so batch load jobs cannot consume the interactive reservation's baseline or autoscaled slots. Assigning interactive queries higher priority within their own reservation guarantees they are scheduled ahead of batch work.

Why this answer

BigQuery reservations allow you to create separate reservations for different workloads (e.g., interactive queries vs. batch loads) and assign them different priority levels. By creating a dedicated reservation for interactive queries with a higher priority, you ensure that interactive queries get preferential access to slots, even when autoscaling slots are consumed by long-running batch jobs. This directly addresses the contention issue without reducing overall capacity.

Exam trap

Google often tests the misconception that autoscaling alone or reducing baseline slots can solve priority issues, but the key is that without separate reservations and explicit priority assignments, all jobs compete equally for the same pool of slots.

How to eliminate wrong answers

Option B is wrong because reducing baseline slots to 200 and relying solely on autoscaling does not solve the priority issue; autoscaling slots are shared and batch jobs could still consume them, leading to the same throttling of interactive queries. Option C is wrong because switching to on-demand pricing eliminates slot reservations entirely, meaning you lose the ability to guarantee capacity or prioritize workloads, and you may face unpredictable performance and higher costs. Option D is wrong because setting the autoscaling max to 1000 slots for batch jobs does not prevent batch jobs from consuming all available slots; it only limits the maximum they can use, but without priority assignment, interactive queries can still be throttled if batch jobs fill the reservation.

15
MCQmedium

Your streaming Dataflow pipeline reads from Pub/Sub, enriches data with a side input, and writes to BigQuery. You need to update the enrichment logic without draining the pipeline, to minimize data loss and maintain exactly-once semantics. What should you do?

A.Cancel the pipeline and create a new one with the updated code.
B.Stop the pipeline, update the code, and restart from the latest snapshot.
C.Use the Dataflow job update mechanism to replace the pipeline with a new version.
D.Drain the pipeline, update the code, and restart with the same job ID.
AnswerC

Dataflow allows updating a streaming pipeline with a new job graph, preserving state and exactly-once processing.

Why this answer

The Dataflow job update mechanism allows you to replace a running pipeline's code with a new version without draining or stopping it, preserving the existing state and minimizing data loss. This mechanism supports exactly-once semantics by ensuring that all in-flight elements are processed exactly once, even after the update, by maintaining the pipeline's checkpoint and watermark state.

Exam trap

The trap here is that candidates often confuse the Dataflow job update mechanism with draining or snapshot-based restarts, not realizing that Dataflow's update feature is specifically designed to allow in-place code changes without data loss or reprocessing.

How to eliminate wrong answers

Option A is wrong because canceling the pipeline would discard all in-flight data and state, leading to data loss and violating exactly-once semantics. Option B is wrong because stopping the pipeline and restarting from a snapshot is not a supported operation in Dataflow; snapshots are used for draining or saving state, but restarting from a snapshot does not guarantee exactly-once processing and can cause data duplication or loss. Option D is wrong because draining the pipeline would allow it to finish processing all existing data before stopping, but then you must create a new pipeline with a new job ID; restarting with the same job ID is not possible after draining, and the drain process itself can cause data loss if not handled correctly.

16
MCQeasy

A data engineer needs to run a recurring SQL transformation in BigQuery every night at 02:00 and, if it fails, retry automatically and send a notification. The team wants the least operational overhead and no external orchestrator. What should they use?

A.A Cloud Composer DAG that runs a BigQueryInsertJobOperator on a cron schedule.
B.A Cloud Scheduler job that calls the BigQuery jobs.insert API with a service account.
C.A BigQuery scheduled query configured with a schedule, a destination table, and notification settings.
D.A Dataflow batch pipeline that reads the source table and writes the transformed result.
AnswerC

BigQuery scheduled queries natively support a recurring schedule, a SQL statement, a destination table, and optional Pub/Sub notification on failure. They run entirely inside BigQuery with no cluster or orchestrator to manage, which matches the requirement for minimal operational overhead. This is the purpose-built feature for recurring SQL transformations on a schedule.

Why this answer

BigQuery scheduled queries are the native mechanism for running SQL on a recurring schedule. They accept a schedule, a query, an optional destination table, and configuration for failure notification through Pub/Sub, all without any infrastructure to operate. Composer, Cloud Scheduler with API calls, and Dataflow all can run SQL, but each introduces components to manage, which conflicts with the requirement for the least operational overhead.

Exam trap

The trap here is assuming a general-purpose orchestrator is always the right answer for scheduled work, when BigQuery's built-in scheduled query already covers a single recurring SQL statement.

17
MCQhard

A company runs a Dataflow streaming pipeline that processes financial transactions. They need to apply a new transformation that enriches the data with a lookup from Cloud Bigtable without stopping the pipeline. The pipeline must be updated in a way that minimises data loss and preserves exactly-once semantics. What is the recommended approach?

A.Use the Dataflow update option with the same pipeline name and new version, ensuring the transform is backward compatible.
B.Drain the pipeline first, then start a new pipeline with the updated code.
C.Create a new pipeline in parallel and switch the Pub/Sub subscription to the new pipeline.
D.Stop the pipeline, update the code, and restart with a new pipeline name.
AnswerA

Dataflow's update option replaces the pipeline definition while retaining the same name and state, so the Bigtable enrichment transform is applied without draining the pipeline, preserving exactly-once semantics provided the new transform stays backward compatible.

Why this answer

Dataflow supports in-place pipeline updates via the update option, which preserves the pipeline's state (including watermarks and deduplication state) and maintains exactly-once semantics. The new transform must be backward compatible with the existing pipeline's state and schema so that the update can be applied without draining.

Exam trap

PDE often tests the difference between update, drain, and stop — the trap is choosing drain or a parallel pipeline when the question requires preserving exactly-once semantics and minimizing data loss.

How to eliminate wrong answers

Option B is wrong because draining stops ingestion and waits for in-flight data to finish, causing downtime and potentially losing data that arrives during the drain window. Option C is wrong because running a parallel pipeline and switching subscriptions risks duplicate processing and breaks exactly-once guarantees across the cutover. Option D is wrong because stopping and restarting with a new pipeline name discards state, loses exactly-once semantics, and may drop or duplicate in-flight data.

18
MCQmedium

Your organization has a BigQuery flat-rate reservation with 500 slots. During peak hours, queries are queued and you need additional capacity temporarily. You want to add slots for a burst of activity without committing to a long-term purchase. What should you do?

A.Switch to on-demand pricing
B.Use flex slots
C.Create a secondary reservation with autoscaling
D.Purchase additional committed use reservations
AnswerB

Flex slots add short-term BigQuery capacity in 500-slot increments for as little as 60 seconds, then delete automatically. They satisfy the stem's temporary burst requirement without a long-term commitment, unlike annual or three-year commitments, and supplement the existing flat-rate reservation's 500 slots during peak queueing.

Why this answer

Flex slots are a BigQuery feature that allows you to purchase slots for a short, fixed duration (e.g., 1 hour, 1 day, 1 week) without a long-term commitment. They are designed exactly for temporary bursts of query activity, such as peak-hour spikes, and can be added to an existing flat-rate reservation. By using flex slots, you temporarily increase the slot capacity of your reservation, reducing queuing, and then the slots automatically expire after the chosen period, avoiding ongoing costs.

Exam trap

PDE often tests the distinction between temporary and committed capacity options, and candidates may confuse flex slots with autoscaling or committed use discounts, leading them to pick options that imply long-term commitments or automatic scaling rather than a short-term, manual burst.

How to eliminate wrong answers

Option A is wrong because switching to on-demand pricing would remove the flat-rate reservation and its predictable cost model, and on-demand pricing is not a temporary add-on; it's a completely different billing model that may not provide the dedicated capacity needed for the burst. Option C is wrong because creating a secondary reservation with autoscaling is not a temporary solution—autoscaling reservations require a baseline commitment and are meant for sustained, variable workloads, not short bursts; moreover, you cannot simply create a secondary reservation without additional cost commitments. Option D is wrong because purchasing additional committed use reservations involves a long-term commitment (1 or 3 years), which contradicts the requirement to avoid a long-term purchase.

19
MCQmedium

A team wants to enforce data quality rules on BigQuery tables using Dataplex. They need to run column-level checks for null values and row-level checks for value ranges on a schedule. Which Dataplex feature should they use?

A.Dataplex Data Profiling
B.BigQuery stored procedures with scheduled queries
C.Dataplex Data Quality Tasks
D.Cloud DLP inspection jobs
AnswerC

Dataplex Data Quality Tasks run scheduled, rule-based checks directly against BigQuery tables, supporting both column-level null validation and row-level range conditions. This satisfies the stem's requirement for automated, recurring enforcement of data quality rules without external tooling, unlike profiling or discovery features that only observe metadata.

Why this answer

Dataplex Data Quality Tasks allow you to define and run data quality rules on BigQuery tables, including column-level checks (e.g., null checks) and row-level checks (e.g., value ranges). These tasks can be scheduled to run periodically, and they generate results that can be monitored. This is the native Dataplex feature for enforcing data quality.

Exam trap

PDE often tests the distinction between Dataplex Data Quality and Data Profiling, so candidates might choose Data Profiling for rule enforcement when it's actually for analysis.

How to eliminate wrong answers

Option A is wrong because Data Profiling analyzes data to discover statistics and patterns, but does not enforce rules or trigger alerts on violations. Option B is wrong because while stored procedures with scheduled queries can implement checks, they are not a Dataplex feature and require custom coding. Option D is wrong because Cloud DLP is for sensitive data discovery and de-identification, not for general data quality rules like null checks or value ranges.

20
Multi-Selectmedium

A data engineer needs to monitor a Pub/Sub-based streaming pipeline. Which two Cloud Monitoring metrics should be used to detect a backlog of unprocessed messages? (Choose two.)

Select 2 answers
A.subscription/oldest_unacked_message_age
B.topic/byte_cost
C.subscription/num_undelivered_messages
D.topic/send_request_count
E.subscription/ack_message_count
AnswersA, C

oldest_unacked_message_age reports how long the oldest unacknowledged message has waited, directly exposing a growing backlog. Rising age means subscribers are not keeping pace with publishing, which is precisely the unprocessed-message condition the engineer must detect.

Why this answer

The 'subscription/num_undelivered_messages' metric shows the number of messages not yet acknowledged, and 'subscription/oldest_unacked_message_age' indicates how long messages have been waiting. Both help detect backlog.

21
Multi-Selecthard

You manage several Cloud Composer 2 environments that run production DAGs. You must define an alerting strategy that detects when a DAG run fails and when a task is stuck retrying for an unusually long time, using Cloud Monitoring. (Choose two.)

Select 2 answers
A.Create an uptime check against the Airflow web server URL and alert when the HTTP response is not 200.
B.Enable Cloud Audit Logs for the Composer API and alert when any environment update method is called.
C.Create a Monitoring alerting policy on the composer.googleapis.com/environment/healthy metric and notify when it drops below the threshold.
D.Create a Monitoring alerting policy on the airflow task instance duration or retry-related metric exposed through Cloud Monitoring for the environment.
E.Create a log-based alerting policy on the Composer airflow logs that matches the DAG run failure log entry and notifies the on-call channel.
AnswersD, E

Composer exposes Airflow metrics such as task instance duration and retry counts to Cloud Monitoring. An alerting policy on those metrics can detect a task whose runtime or retry behaviour exceeds a normal threshold, which is exactly the stuck-retrying condition described. It complements the failure alert by catching slow degradation before hard failure.

Why this answer

Composer surfaces Airflow execution data in two complementary ways: structured log entries in Cloud Logging and Airflow metrics in Cloud Monitoring. A log-based alerting policy catches the explicit DAG run failure event, while a metric-based policy on task duration or retries catches a task that is hanging or looping. Together they cover both the hard failure and the stuck-retrying condition, whereas health, uptime, and audit signals describe infrastructure and control-plane state rather than workload outcomes.

Exam trap

The trap here is reaching for environment health or web server uptime metrics, which measure whether Composer is running rather than whether your DAGs are succeeding.

22
MCQmedium

Your team runs a Cloud Composer 2 environment (composer-2.1.0-airflow-2.6.3) that executes a daily BigQuery ETL workflow. The workflow must not run on weekends. You want to implement this with minimal code and without modifying the DAG's task logic. What should you do?

A.Add a BranchPythonOperator that checks the day of the week and skips downstream tasks on weekends.
B.Set the DAG's schedule_interval to '@daily' and add a ShortCircuitOperator that stops execution on weekends.
C.Create a Cloud Scheduler job that triggers the Composer DAG only on weekdays via the Airflow REST API.
D.Set the DAG's schedule_interval to '0 2 * * 1-5' and set catchup=False.
AnswerD

The cron expression '0 2 * * 1-5' runs the DAG at 02:00 Monday through Friday, which excludes weekends. Setting catchup=False prevents backfilling missed runs. This directly satisfies the requirement to skip weekends without altering task logic, and it uses the native Airflow scheduling mechanism supported in Cloud Composer.

Why this answer

The correct approach is to use a cron expression that runs only Monday through Friday, such as '0 2 * * 1-5', combined with catchup=False to avoid backfilling. This leverages Airflow's native scheduling, requires no changes to task logic, and keeps the DAG simple. Other methods either add external dependencies or still trigger the DAG on weekends.

Exam trap

The trap here is assuming that a BranchPythonOperator or ShortCircuitOperator prevents the DAG from running on weekends, when in fact it only prevents downstream tasks from executing after the DAG has already been triggered.

23
MCQhard

You manage a Cloud Dataflow streaming pipeline that reads from Pub/Sub and writes to BigQuery. The pipeline uses the BigQueryIO write transform with STREAMING_INSERTS. You need to ensure exactly-once processing semantics for the BigQuery writes. What should you do?

A.Add a BigQuery insertId to each record and rely on BigQuery's built-in deduplication for streaming inserts.
B.Use BigQueryIO.Write with withMethod(STORAGE_WRITE_API) and set withNumStorageWriteApiStreams to a value greater than zero.
C.Switch to BigQueryIO.Write with withMethod(FILE_LOADS) and set withTriggeringFrequency to a low value.
D.Enable exactly-once by setting the pipeline option --experiments=use_runner_v2 and using BigQueryIO.Write with withMethod(STREAMING_INSERTS).
AnswerB

The Storage Write API supports exactly-once semantics when used with Dataflow. Configuring withMethod(STORAGE_WRITE_API) and specifying one or more streams enables the necessary deduplication and transactional guarantees. This is the recommended approach for exactly-once streaming writes to BigQuery, replacing the older STREAMING_INSERTS method.

Why this answer

To achieve exactly-once processing for BigQuery writes in a Dataflow streaming pipeline, use the BigQuery Storage Write API with withMethod(STORAGE_WRITE_API) and configure at least one stream. This method provides transactional writes and deduplication, ensuring each record is written exactly once. Other methods like STREAMING_INSERTS with insertId offer only best-effort deduplication and cannot guarantee exactly-once semantics.

Exam trap

The trap here is assuming that the insertId field in streaming inserts guarantees exactly-once semantics, when in reality it only provides best-effort deduplication within a limited time window.

24
MCQeasy

A data engineer schedules a Cloud Composer 2 environment to run a DAG that triggers a Dataflow batch job every night. The DAG sometimes fails because the Dataflow job takes longer than the default task timeout. The engineer wants the DAG to wait for the Dataflow job to finish rather than timing out. Which change should the engineer make?

A.Add a retry policy to the DAG so failed tasks are retried until the Dataflow job finishes.
B.Increase the DAG's schedule interval so runs start less frequently.
C.Set the Dataflow operator's execution_timeout and deferrable parameters so the task waits for job completion.
D.Move the Dataflow launch into a Cloud Function and call it from the DAG.
AnswerC

The Dataflow operators in Cloud Composer can wait for the job to complete, and setting an appropriate execution_timeout prevents the task from failing early. Using the deferrable mode releases the worker slot while waiting, which is the recommended pattern for long-running jobs in Composer 2. This directly addresses the timeout problem.

Why this answer

Cloud Composer 2 Dataflow operators can block until the job reaches a terminal state, and the deferrable variant frees the worker while waiting. Setting execution_timeout to a value longer than the expected job duration prevents premature task failure. Schedule interval, Cloud Function wrapping, and retry policies do not make a task wait for the underlying Dataflow job.

Exam trap

The trap here is confusing task retries or schedule changes with making the operator wait for an asynchronous job to finish.

25
MCQmedium

A financial services firm stores customer transaction data in BigQuery. Compliance requires that a nightly Cloud Composer DAG verify that the previous day's partition is complete before downstream reporting DAGs run, and that the reporting DAG never start if the verification fails. The engineer wants the dependency expressed inside orchestration rather than by polling from the reporting DAG. What should the engineer do?

A.Create a Cloud Scheduler job that calls the reporting DAG's trigger endpoint only after a Cloud Function confirms verification succeeded.
B.Combine the verification and reporting logic into a single DAG so that task order enforces the dependency.
C.Have the reporting DAG poll the verification table in a loop with a sensor until a success row appears, then proceed.
D.Use a Dataset (Dataset-scheduled DAG) in Cloud Composer: have the verification DAG produce the dataset and schedule the reporting DAG to trigger on it.
AnswerD

Airflow datasets provide a push-based dependency: the verification DAG marks the dataset as updated on success, and the reporting DAG is scheduled to run only when that dataset is updated. This expresses the dependency declaratively in orchestration, satisfies the no-polling requirement, and naturally blocks reporting when verification fails.

Why this answer

Airflow datasets in Cloud Composer create a producer-consumer dependency where the verification DAG's successful completion updates a dataset and the reporting DAG is scheduled on that update. This keeps the dependency inside orchestration, avoids polling, and inherently prevents the reporting DAG from running when verification fails.

Exam trap

The trap here is defaulting to a sensor or an external scheduler for cross-DAG dependencies, when Airflow's dataset scheduling is the push-based mechanism designed for exactly this.

26
MCQhard

A financial services company runs a Dataflow streaming pipeline that reads from Pub/Sub and writes enriched records to BigQuery. Compliance requires that the raw Pub/Sub messages be retained for seven years so that any record can be reprocessed if the enrichment logic is later found to be incorrect. The pipeline currently has no archival step. What should the data engineer do to satisfy the retention requirement with the least operational overhead?

A.Configure a Pub/Sub dead-letter topic and set the maximum delivery attempts high enough that failed messages persist for the retention period.
B.Enable BigQuery table snapshots on the destination table and set a seven-year expiration on each snapshot.
C.Add a Cloud Storage sink to the Dataflow pipeline that writes each raw message to a Cloud Storage bucket with a lifecycle policy moving objects to Coldline or Archive storage after 30 days.
D.Increase the Pub/Sub topic message retention duration to the maximum allowed and rely on the subscription backlog for seven years.
AnswerC

Writing the unmodified messages to Cloud Storage preserves the raw payload independently of the BigQuery output, and a lifecycle policy automatically transitions old objects to colder, cheaper classes so seven years of retention stays affordable. The Dataflow sink adds little operational burden because the pipeline already runs continuously and only needs one additional write branch.

Why this answer

Retaining raw messages for years requires a durable, low-cost object store, and Cloud Storage with lifecycle-based class transitions is the standard fit on Google Cloud. By branching the existing Dataflow pipeline to write unmodified messages to a bucket, the company keeps a complete, reprocessable archive without building a separate ingestion path or managing long-lived Pub/Sub backlogs.

Exam trap

The trap here is assuming Pub/Sub retention can be stretched to meet a multi-year compliance window, when its retention is bounded and it is not an archival system.

27
Multi-Selecteasy

You are using Cloud Workflows to orchestrate a series of API calls. You need to handle errors and retries. Which THREE features of Cloud Workflows can you use? (Choose THREE.)

Select 3 answers
A.Use try/except blocks to catch and handle errors.
B.Integrate with Cloud Load Balancing for high availability.
C.Use conditional branches (if-else) based on step results.
D.Define a retry policy on a step.
E.Enable automatic logging for each step.
AnswersA, C, D

Cloud Workflows supports try/except syntax within a step, letting you catch a raised error and route execution to a handler instead of failing the whole execution. This directly satisfies the need to handle errors during orchestration of the API call sequence.

Why this answer

Option A is correct because Cloud Workflows supports try/except blocks, allowing you to catch errors raised by a step and execute fallback logic or handle the failure gracefully. Option C is correct because Cloud Workflows supports conditional branching with if/else constructs, letting you evaluate step results and choose different execution paths based on those outcomes. Option D is correct because Cloud Workflows lets you attach a retry policy to a step, specifying parameters such as max_retries, initial_delay, max_delay, and multiplier to automatically retry failed calls.

Option B is not correct because Cloud Load Balancing is a traffic-distribution service for workloads, not a Cloud Workflows error-handling or retry feature. Option E is not correct because logging is handled by Cloud Logging integration and is not a configurable per-step error-handling or retry feature of Cloud Workflows.

Exam trap

The trap is selecting logging or load balancing as error-handling features; candidates must distinguish between observability (logging) and actual error/retry control mechanisms (try/except, retry policies).

28
MCQhard

Your organization runs a Cloud Composer 2 environment that executes dozens of DAGs. Several DAGs share a connection to an external REST API that enforces a rate limit of 100 requests per minute. During peak hours, DAGs fail with HTTP 429 errors. You want to prevent these failures without changing the external API's limits. What should you do?

A.Increase the number of Celery workers in the Composer environment.
B.Set the DAG's max_active_runs to 1 for each affected DAG.
C.Add exponential backoff retries to the API-calling tasks.
D.Create an Airflow pool with a limited number of slots and assign the affected tasks to that pool.
AnswerD

Airflow pools limit the number of concurrent task instances that can use a set of slots. By creating a pool sized to stay under the API's rate limit and assigning the API-calling tasks to it, you throttle concurrency across all DAGs. This directly prevents exceeding 100 requests per minute without modifying the external service.

Why this answer

Airflow pools provide a shared concurrency limit that spans DAGs. Sizing a pool to keep total in-flight API calls under the provider's rate limit prevents 429 responses. Increasing workers, limiting per-DAG active runs, or relying on retries does not coordinate request volume across the many DAGs that share the external API.

Exam trap

The trap here is treating retries or per-DAG concurrency limits as a substitute for a cross-DAG concurrency control like an Airflow pool.

29
MCQhard

You have a Dataflow batch pipeline that processes data from Cloud Storage and writes to BigQuery. The pipeline uses a custom DoFn that sometimes throws exceptions due to malformed input records. You want to ensure that the pipeline continues processing valid records while logging the malformed ones for later analysis, without failing the entire job. Which Dataflow feature should you use?

A.Configure the pipeline to use a dead-letter queue in Pub/Sub.
B.Use a side output to capture and log the malformed records.
C.Set the pipeline's failure mode to 'continue' in the pipeline options.
D.Use a ParDo with a try-catch block and log the exception, then drop the record.
AnswerB

Dataflow supports side outputs, which allow you to route elements that fail processing to a separate PCollection. By catching exceptions in your DoFn and emitting the malformed records to a side output, you can continue processing valid records without failing the pipeline. The side output can then be written to a dead-letter sink, such as Cloud Storage or BigQuery, for later analysis. This is the recommended pattern for handling bad records in Dataflow.

Why this answer

The correct approach is to use a side output to capture malformed records. Side outputs allow you to emit elements that fail processing to a separate PCollection, which can then be written to a dead-letter sink for analysis. This enables the pipeline to continue processing valid records without failing.

Other options either do not exist, drop data, or are not applicable to batch pipelines. Side outputs are a core Dataflow feature for error handling.

Exam trap

The trap here is assuming that simply logging and dropping bad records is sufficient, when the requirement is to preserve them for later analysis, which necessitates a side output to a separate sink.

30
MCQeasy

A data engineer must ensure that a Cloud Composer DAG which loads a BigQuery table runs every day at 02:00 UTC and that a dependent downstream report DAG runs only after the load succeeds. The report DAG lives in the same Composer environment but is a separate DAG file. Which feature should be used to coordinate the two DAGs?

A.A shared XCom written by the load DAG and read by the report DAG using xcom_pull across DAGs.
B.A single trigger_rule of all_success applied to the report DAG's root task.
C.A TimeSensor in the report DAG set to the expected load completion time.
D.An ExternalTaskSensor in the report DAG referencing the load DAG and task, with a matching execution date.
AnswerD

ExternalTaskSensor is designed to wait on a task in a different DAG, which matches the cross-DAG dependency here. With the correct execution_delta or execution_date_fn so the two runs align, the report DAG waits until the load task succeeds before proceeding, giving the required success-based coordination.

Why this answer

ExternalTaskSensor is the purpose-built mechanism for waiting on a task in another DAG. Configuring it with the correct execution date alignment makes the report DAG block until the load task succeeds, which is exactly the cross-DAG dependency described. Time-based and intra-DAG mechanisms cannot observe another DAG's success state.

Exam trap

The trap here is using a time-based wait or an intra-DAG trigger rule when the dependency crosses DAG boundaries.

31
MCQeasy

An engineer needs to create a reusable Dataflow pipeline that can be executed with different parameters without modifying code. Which Dataflow feature should they use?

A.Dataflow Shuffle
B.Dataflow Flex Templates
C.Dataflow SQL
D.Dataflow Classic Templates
AnswerB

Flex Templates package the pipeline as a Docker image with a metadata specification, so the same template runs repeatedly with different runtime parameters and no code changes. This satisfies the stem's reusability and parameterisation constraint.

Why this answer

Dataflow Flex Templates allow packaging a pipeline as a Docker image with a metadata file, enabling reuse with different parameters at runtime without code changes. They support dynamic parameters and are the recommended approach for reusable pipelines.

Exam trap

PDE often tests the difference between Classic and Flex Templates, where candidates incorrectly choose Classic Templates for dynamic parameterization.

How to eliminate wrong answers

Option A is wrong because Dataflow Shuffle is a service for shuffling data, not for pipeline templating. Option C is wrong because Dataflow SQL is a way to run SQL queries on Dataflow, not for creating reusable parameterized pipelines. Option D is wrong because Classic Templates require parameters to be defined at compile time and are less flexible than Flex Templates.

32
MCQeasy

Your company uses Cloud Composer to run a daily ETL workflow. The workflow consists of several tasks that must run in a specific order. You want to receive an alert if any task fails. Which Cloud Monitoring feature should you use?

A.Create an alerting policy based on the Cloud Composer metric for task failures.
B.Set up a Cloud Function that checks the Airflow UI periodically.
C.Configure Airflow email alerts in the DAG's default_args.
D.Use Cloud Logging to create a log-based metric for task failures.
AnswerA

Cloud Composer exports metrics to Cloud Monitoring, including metrics for task failures. You can create an alerting policy that triggers when the number of failed tasks exceeds a threshold. This is the most direct way to get notified of task failures in your ETL workflow. It leverages the native integration between Cloud Composer and Cloud Monitoring, allowing you to set up notifications via email, SMS, or other channels.

Why this answer

The recommended way to alert on task failures in Cloud Composer is to use Cloud Monitoring alerting policies based on the built-in Cloud Composer metrics for task failures. These metrics are automatically exported and can be used to trigger alerts when failures occur. This approach is native, scalable, and integrates with various notification channels.

Other methods, such as log-based metrics or polling, are either redundant or less efficient.

Exam trap

The trap here is overlooking the built-in Cloud Composer metrics and instead opting for custom log-based metrics or polling, which are unnecessary and more complex.

33
MCQeasy

A streaming Dataflow pipeline needs to be updated without draining the existing pipeline. Which update strategy should be used?

A.Drain the pipeline first, then start a new one
B.Replace the job with a new job using the same pipeline name
C.Use a different pipeline name and cancel the old one
D.Stop the job, update the code, and restart
AnswerB

A replacement job cannot be updated in place, so the pipeline is relaunched under the same name, letting the new job start while the old one is cancelled. This avoids the drain step that would stop ingestion, satisfying the no-drain constraint.

Why this answer

Dataflow supports replacing a running streaming job with an updated one using the same pipeline name via the 'Replace job' (or update) mechanism, which performs an in-place upgrade that preserves the pipeline's state and does not require draining. This is the supported path for updating streaming pipelines without stopping data ingestion. Draining, cancelling, or stopping the job all interrupt processing and are not the recommended no-drain update strategy.

Exam trap

The trap is assuming any update requires stopping or draining the pipeline; candidates pick the 'safe-sounding' drain option without realizing Dataflow's replace-job feature is specifically designed to avoid draining.

How to eliminate wrong answers

Option A is wrong because draining stops the pipeline from ingesting new data and waits for in-flight data to finish, which is the opposite of updating without draining. Option C is wrong because using a different pipeline name and cancelling the old one creates a brand-new job with no state continuity and causes a gap in processing. Option D is wrong because stopping the job, editing code, and restarting is a manual stop/start cycle that halts the pipeline and loses the seamless update semantics Dataflow provides.

34
MCQhard

A data engineer is migrating a Composer 1 environment to Cloud Composer 2 and notices that a DAG relying on a legacy operator for a deprecated service no longer works. The engineer wants a durable fix that keeps the pipeline running and avoids repeating this problem in future upgrades. Which action should be taken?

A.Downgrade the environment's Python version to match the requirements of the legacy operator.
B.Move the deprecated logic into a PythonOperator that calls the service's current client library, and add unit tests to validate the behaviour.
C.Pin the Composer 2 environment to the exact Airflow version that still includes the legacy operator and avoid upgrades.
D.Copy the legacy operator's source into the DAGs folder so it is imported locally instead of from the provider package.
AnswerB

Replacing the removed operator with a PythonOperator that uses the service's supported client library removes the dependency on deprecated code and keeps the DAG working. Adding unit tests guards against regressions during future upgrades, giving a durable, maintainable solution rather than a temporary workaround.

Why this answer

Replacing deprecated operators with a PythonOperator that calls the service's current client library removes reliance on unsupported code and keeps the DAG functional on Composer 2. Adding tests ensures the replacement behaves correctly and catches regressions during future upgrades, which is the durable, maintainable outcome the scenario asks for.

Exam trap

The trap here is treating version pinning or vendoring deprecated code as a fix, when it only postpones the incompatibility.

35
MCQeasy

Which BigQuery feature allows you to estimate the cost of a query before running it, by returning the number of bytes that would be processed?

A.EXPLAIN statement
B.INFORMATION_SCHEMA.JOBS
C.--dry_run flag
D.Slot estimator
AnswerC

The `--dry_run` flag validates a query and returns the bytes it would process without executing it, so no compute charges are incurred. This directly satisfies the stem's requirement to estimate cost before running, since BigQuery on-demand pricing is calculated from bytes processed.

Why this answer

The --dry_run flag in BigQuery allows you to validate a query and estimate the number of bytes it will process without actually running it, thus providing a cost estimate. This is a standard feature in the bq command-line tool and the BigQuery API. It returns the total bytes processed, which can be used to calculate the cost based on current pricing.

Exam trap

The trap is confusing the dry run flag with other features like EXPLAIN or INFORMATION_SCHEMA, which serve different purposes.

How to eliminate wrong answers

Option A is wrong because EXPLAIN is not a BigQuery feature; it is used in other databases to show query plans. Option B is wrong because INFORMATION_SCHEMA.JOBS provides metadata about completed jobs, not a pre-execution cost estimate. Option D is wrong because the slot estimator is not a feature for estimating query cost; slots are compute resources, and there is no direct 'slot estimator' for cost estimation.

36
MCQeasy

A data engineer has a BigQuery SQL script that must run every day at 06:00, load its results into a reporting table, and retry automatically if the query fails due to transient errors. The team has no existing orchestration tooling, wants the lowest operational overhead, and needs the schedule and the SQL to be managed entirely inside Google Cloud. Which approach should the engineer use?

A.Use Dataflow with a batch pipeline that reads from the source tables, applies the transformation, and writes to the reporting table on a schedule.
B.Create a Cloud Composer environment and author a single-task DAG that runs the SQL daily with Airflow retries configured.
C.Create a scheduled query in BigQuery with the desired SQL and a daily schedule, targeting the reporting table as the destination.
D.Deploy the SQL as a Cloud Run service and create a Cloud Scheduler HTTP job that invokes it daily at 06:00.
AnswerC

BigQuery scheduled queries let you store SQL, define a schedule, and write results to a destination table directly in the BigQuery UI or API. The service manages execution and retries failed runs automatically, requiring no containers, schedulers, or external orchestration, which is exactly the lowest-overhead option described.

Why this answer

BigQuery scheduled queries are the native, serverless way to run SQL on a recurring schedule and write output to a destination table, with automatic retry of failed runs. Because the schedule, SQL, and destination all live inside BigQuery, no additional infrastructure or orchestration tooling is needed, matching the low-overhead requirement precisely.

Exam trap

The trap here is reaching for external schedulers or orchestration platforms when BigQuery's own scheduled query feature already covers recurring SQL with retries and destination tables.

37
MCQmedium

A data engineer is building a batch pipeline that runs daily using Cloud Composer. The pipeline has three tasks: extract data from Cloud Storage, transform data using Dataflow, and load the transformed data into BigQuery. The engineer wants to ensure that the Dataflow job only starts after the extraction task completes successfully, and the load task only starts after the Dataflow job finishes. How should the engineer define the task dependencies in the Airflow DAG?

A.extract >> [transform, load]
B.transform >> extract >> load
C.extract >> transform >> load
D.extract >> load >> transform
AnswerC

Airflow's bitshift operator sets explicit downstream dependencies, so transform runs only after extract succeeds and load only after transform completes. This directly enforces the sequential ordering the stem requires, preventing the Dataflow job from starting prematurely.

Why this answer

Airflow uses the bitshift operator (>>) to define task dependencies in a linear sequence. The DAG must ensure that the extract task completes before the transform task starts, and the transform task completes before the load task starts. This is achieved by chaining the tasks in order: extract >> transform >> load, which enforces the required sequential execution.

Exam trap

A common misconception is that multiple tasks can be chained in parallel with a single bitshift operator, leading candidates to choose Option A, which incorrectly allows the load task to start before the Dataflow job completes. In Airflow, sequential dependencies are defined by chaining tasks with '>>' in order.

How to eliminate wrong answers

Option A is wrong because it sets transform and load as parallel downstream tasks of extract, meaning load could start before transform finishes, violating the requirement that load waits for Dataflow. Option B is wrong because it places transform before extract, which would attempt to run the Dataflow job before the extraction completes, breaking the dependency chain. Option D is wrong because it places load before transform, which would attempt to load data into BigQuery before the Dataflow transformation is done, leading to incorrect or missing data.

38
MCQeasy

A data engineer wants to quickly estimate the cost of running a BigQuery query before executing it. Which command-line tool or command should they use?

A.gcloud logging read
B.bq query --use_cache=false
C.gcloud bigtable queries run
D.bq query --dry_run
AnswerD

The bq query --dry_run flag validates the query and returns the bytes it would process without executing it or incurring charges. Multiplying those bytes by the on-demand rate gives a cost estimate before running the job.

Why this answer

The bq query --dry_run flag validates a query and returns the amount of data it would process without actually executing it, which is the standard way to estimate BigQuery on-demand query cost. Since BigQuery on-demand pricing is based on bytes scanned (currently $6.25 per TiB), the dry-run's reported bytes give a direct cost estimate. It also validates syntax and permissions without incurring charges.

Exam trap

PDE often tests the distinction between flags that change execution behavior (--use_cache=false) and flags that only estimate without executing (--dry_run), so candidates who don't know the dry-run semantics pick a plausible-sounding but wrong flag.

How to eliminate wrong answers

Option A is wrong because gcloud logging read retrieves log entries and has nothing to do with query cost estimation. Option B is wrong because --use_cache=false disables result caching, which would if anything increase cost by forcing a full scan — it does not estimate anything. Option C is wrong because gcloud bigtable queries run is not a valid BigQuery cost tool and targets Bigtable, a different NoSQL service entirely.

39
MCQeasy

You need to schedule a BigQuery query to run every day at 6:00 AM and write the results to a new table. The query is simple and does not require complex dependencies. You want a low-maintenance, serverless solution with minimal configuration. What should you do?

A.Use BigQuery scheduled queries to run the query daily and write to a destination table.
B.Deploy a Cloud Function triggered by Cloud Scheduler to execute the query via the BigQuery API.
C.Create a Dataflow pipeline with a scheduled trigger to run the query and write results.
D.Create a Cloud Composer DAG with a BigQueryInsertJobOperator scheduled with a cron expression.
AnswerA

BigQuery scheduled queries are a native, serverless feature that allows you to schedule SQL queries with a specified frequency. They require minimal configuration and no infrastructure management. You can set the schedule, destination table, and even notifications. This is the simplest solution for a single daily query.

Why this answer

BigQuery scheduled queries are a built-in, serverless feature that lets you schedule SQL queries to run at specified intervals. They require no infrastructure management and can write results to a destination table. For a simple daily query, this is the most straightforward and low-maintenance solution, avoiding the overhead of Cloud Composer, Cloud Functions, or Dataflow.

Exam trap

The trap here is over-engineering a simple scheduling requirement by choosing orchestration tools like Cloud Composer or custom code with Cloud Functions, when BigQuery's native scheduling suffices.

40
MCQeasy

You want to monitor the latency of messages in a Pub/Sub subscription. Which Cloud Monitoring metric should you use to see the age of the oldest unacknowledged message?

A.pubsub.googleapis.com/subscription/oldest_unacked_message_age
B.pubsub.googleapis.com/subscription/num_undelivered_messages
C.pubsub.googleapis.com/topic/send_request_count
D.pubsub.googleapis.com/topic/publish_latency
AnswerA

The metric `pubsub.googleapis.com/subscription/oldest_unacked_message_age` directly reports the age of the oldest unacknowledged message in a subscription, measured in seconds. This precisely satisfies the stem's requirement to monitor message latency, since backlog age reflects how long messages wait before acknowledgement.

Why this answer

The Cloud Monitoring metric pubsub.googleapis.com/subscription/oldest_unacked_message_age reports the age (in seconds) of the oldest unacknowledged message in a subscription, which is exactly the latency indicator requested. It is a subscription-level metric that directly reflects how long a message has been waiting without acknowledgment, making it the correct choice for monitoring message age.

Exam trap

PDE often tests whether candidates confuse publisher-side latency metrics (publish_latency, send_request_count) with subscription-side backlog/age metrics (oldest_unacked_message_age, num_undelivered_messages).

How to eliminate wrong answers

Option B is wrong because num_undelivered_messages counts how many messages are pending, not how old the oldest one is, so it measures backlog size rather than latency. Option C is wrong because topic/send_request_count measures publish request volume at the topic level, not message age in a subscription. Option D is wrong because topic/publish_latency measures the time to publish to the topic, which is a publisher-side latency, not the age of unacknowledged messages in a subscription.

41
MCQmedium

You are building a data pipeline that runs daily batch jobs on Dataproc, then loads results into BigQuery. You want to orchestrate the entire workflow, including dependencies between steps, retries, and monitoring. Which Google Cloud service is most appropriate?

A.Cloud Scheduler
B.Cloud Composer
C.Cloud Workflows
D.Dataflow
AnswerB

Cloud Composer, built on Apache Airflow, models the pipeline as a directed acyclic graph, so Dataproc job dependencies, retries and monitoring are orchestrated natively. It satisfies the stem's requirement for cross-service workflow orchestration, which BigQuery scheduled queries alone cannot provide.

Why this answer

Cloud Composer is the most appropriate service because it is a fully managed Apache Airflow service that provides workflow orchestration with DAGs, dependency management, retries, scheduling, and monitoring. It natively integrates with Dataproc and BigQuery operators, allowing you to define the entire pipeline as code. This matches the requirement for orchestrating a multi-step daily batch workflow.

Exam trap

The trap here is confusing orchestration services: candidates often pick Cloud Workflows because it sounds like a workflow tool, but the exam expects you to recognize that Airflow-based Cloud Composer is the standard for data pipeline orchestration with dependencies and retries.

How to eliminate wrong answers

Option A is wrong because Cloud Scheduler is only a cron-based job scheduler that triggers a single target (HTTP, Pub/Sub, App Engine); it cannot manage dependencies, retries, or multi-step workflows. Option C is wrong because Cloud Workflows is a serverless orchestration engine for API calls and service chaining, but it lacks the rich operator ecosystem and DAG-based dependency management of Airflow, making it less suited for data pipeline orchestration. Option D is wrong because Dataflow is a data processing service for stream and batch pipelines, not a workflow orchestrator; it does not manage dependencies between Dataproc and BigQuery steps.

42
MCQhard

A Dataflow streaming pipeline writes to BigQuery and has run in production for months. The team wants to add a transformation and deploy the change with zero data loss and no interruption to the running pipeline. Which deployment approach should they use?

A.Run the updated pipeline in parallel with the old one, then delete the old pipeline after a fixed time.
B.Use a Cloud Scheduler job to redeploy the pipeline template on a cron and let the new job take over.
C.Update the pipeline in place by calling the Dataflow update method with the new job graph and a compatible transform name.
D.Stop the existing pipeline, then start a new pipeline with the updated code and the same job name.
AnswerC

Dataflow supports updating a running streaming job with a new pipeline definition when the transforms are named consistently and the update is compatible, preserving the existing state such as windows and timers. The service swaps in the new graph while maintaining exactly-once semantics for the sinks, so processing continues without draining or losing in-flight data. This is the intended zero-downtime deployment path.

Why this answer

Dataflow's update capability is the supported way to change a running streaming pipeline without stopping it. When the new graph keeps transform names consistent and the changes are compatible, the service swaps the graph while preserving streaming state and maintaining the sink's exactly-once guarantees. Stopping and restarting, running duplicate pipelines, or redeploying templates all either create a processing gap, risk data loss, or cause duplicate output.

Exam trap

The trap here is treating a pipeline update like a code redeploy, where any restart is fine, instead of recognizing that streaming state and exactly-once sinks must be preserved across the change.

43
MCQeasy

Which Dataflow feature allows you to package a pipeline into a reusable template that can be deployed with different parameters at runtime?

A.Cloud Dataproc
B.Classic Templates
C.Dataflow SQL
D.Flex Templates
AnswerD

Flex Templates package the pipeline as a Docker image plus a metadata file, allowing runtime parameterisation when launched via the gcloud CLI, REST API or Cloud Scheduler. This satisfies the requirement to reuse one pipeline definition across differing runtime parameters.

Why this answer

Dataflow Flex Templates allow you to containerize a pipeline and provide runtime parameters, enabling reusability across different environments or jobs.

44
MCQmedium

You need to inspect a BigQuery table for sensitive data such as credit card numbers and apply masking. Which GCP service should you use to identify and de-identify the data?

A.IAM Recommender
B.Cloud KMS
C.Dataplex
D.Cloud Data Loss Prevention (DLP)
AnswerD

Cloud Data Loss Prevention scans BigQuery tables using built-in infoTypes to detect credit card numbers, then applies de-identification transforms such as masking, tokenisation or redaction directly on the findings. This satisfies the stem's requirement to both identify sensitive data and mask it within BigQuery.

Why this answer

Cloud Data Loss Prevention (DLP) is Google Cloud's service for discovering, classifying, and de-identifying sensitive data such as credit card numbers, SSNs, and emails. It provides infoType detectors (including CREDIT_CARD_NUMBER) and de-identification transforms like masking, tokenization, and format-preserving encryption. It integrates directly with BigQuery for scanning tables and applying masking at query or storage time.

Exam trap

PDE often tests the confusion between governance/orchestration services (Dataplex) and the actual detection-and-de-identification engine (Cloud DLP), so candidates pick Dataplex when the task is specifically identifying and masking sensitive data.

How to eliminate wrong answers

Option A is wrong because IAM Recommender suggests least-privilege role changes based on usage — it has no data inspection or de-identification capability. Option B is wrong because Cloud KMS manages encryption keys for data at rest; it encrypts but does not identify sensitive content or apply field-level masking. Option C is wrong because Dataplex is a data governance and lakehouse management service that can orchestrate discovery and lineage, but the actual sensitive-data detection and de-identification engine it invokes is Cloud DLP, not Dataplex itself.

45
MCQmedium

Your team uses Cloud Composer to orchestrate a daily ETL workflow that extracts data from an on-premises database, transforms it in Dataproc, and loads it into BigQuery. The workflow occasionally fails due to transient network issues. You want to automatically retry failed tasks without manual intervention. What should you do?

A.Configure the Dataproc cluster to automatically restart on failure.
B.Enable automatic retries on the Cloud Composer environment.
C.Use Cloud Scheduler to trigger the DAG multiple times until it succeeds.
D.Set the retries parameter and retry_delay on the Airflow tasks in the DAG.
AnswerD

In Cloud Composer, Airflow tasks support retries and retry_delay parameters. By setting these, you can automatically retry failed tasks after a specified delay. This is the standard way to handle transient failures in Airflow DAGs. You can also set exponential backoff for retries. This approach requires no additional infrastructure and is fully integrated with Cloud Composer.

Why this answer

Airflow tasks in Cloud Composer support retries and retry_delay parameters, which allow automatic retries of failed tasks. This is the correct way to handle transient failures. Other options either do not address task-level retries or could cause data duplication.

Exam trap

The trap here is thinking that Cloud Composer has an environment-wide retry setting, but retries must be configured per task.

46
Multi-Selectmedium

A team runs a Dataflow batch pipeline that writes to BigQuery. They want to reduce cost and improve reliability of the write step. The pipeline currently uses BigQueryIO with the default write method. Which two changes should they make to achieve these goals? (Choose two.)

Select 2 answers
A.Increase the number of BigQuery load jobs by setting a very low triggering frequency so each micro-batch is written separately.
B.Set the BigQueryIO write disposition to WRITE_TRUNCATE on every micro-batch so the table is always consistent.
C.Enable the withAutoSharding option on the write transform so the number of shards adapts to the pipeline's throughput.
D.Disable dynamic work rebalancing so that workers process a fixed set of shards for the entire pipeline run.
E.Switch the write to the Storage Write API with exactly-once semantics instead of the legacy streaming inserts or load jobs.
AnswersC, E

withAutoSharding dynamically adjusts the number of shards used for the BigQuery write based on the pipeline's throughput, which improves performance and reduces the chance of a single shard becoming a bottleneck. It is compatible with the Storage Write API and helps the write step scale without manual tuning. This supports reliability and cost efficiency by avoiding over-provisioning of shards for low-volume periods.

Why this answer

The Storage Write API lowers ingestion cost and supports exactly-once semantics, addressing both cost and reliability. Enabling withAutoSharding lets the write step adapt shard count to actual throughput, preventing shard bottlenecks and avoiding over-provisioning. Together these changes reduce per-row cost and improve the robustness of the BigQuery write step without altering the pipeline's business logic.

Exam trap

The trap here is assuming that more frequent, smaller load jobs improve reliability, when they actually increase quota pressure and metadata overhead while raising cost.

47
Multi-Selectmedium

A company is migrating their on-premises data warehouse to BigQuery. They have a mix of batch and streaming ingestion. The data team wants to optimize query costs. Which THREE practices should they adopt?

Select 3 answers
A.Switch to flat-rate pricing to cap slot usage.
B.Use materialized views for frequently executed aggregations.
C.Partition tables by a date or timestamp column.
D.Limit the number of concurrent queries by setting a maximum slot capacity.
E.Cluster tables on columns that are frequently used in filters and joins.
AnswersB, C, E

Materialised views precompute and store aggregation results, so repeated queries scan the small stored result rather than the full base tables. This directly satisfies the cost-optimisation goal by reducing bytes processed for frequently executed aggregations.

Why this answer

Option B is correct because materialized views precompute and store the results of frequently executed aggregations, so repeated queries read the smaller cached result set instead of rescanning the full base tables, directly reducing bytes processed and thus query cost. Option C is correct because partitioning a table by a DATE or TIMESTAMP column enables partition pruning, so queries with filters on that column scan only the relevant partitions rather than the entire table, cutting the data billed per query. Option E is correct because clustering on columns frequently used in filters and joins co-locates related data in storage blocks, allowing BigQuery to skip irrelevant blocks and reduce bytes scanned for those predicates.

Option A is not correct because flat-rate pricing changes how slots are billed (a committed capacity model) rather than reducing the bytes processed by queries, so it does not itself optimize query cost for this migration. Option D is not correct because capping concurrent queries or slot capacity is a workload-management throttle, not a cost-optimization practice, and it can even slow throughput without lowering bytes-scanned charges.

Exam trap

PDE often tests the confusion between cost optimization and performance/capacity management; candidates pick flat-rate or concurrency limits thinking they reduce cost, but only reducing bytes scanned via partitioning, clustering, and materialized views lowers on-demand query costs.

48
Multi-Selectmedium

You want to optimize BigQuery costs for a large dataset that is frequently queried by time range. You also need to ensure that predictable workloads have dedicated slot capacity. Which TWO strategies should you combine? (Choose 2)

Select 2 answers
A.Use query caching
B.Partition the table by date
C.Purchase committed use reservations for baseline capacity
D.Create a materialized view for the entire table
E.Enable autoscaling slots
AnswersB, C

Partitioning by date restricts each query to the relevant date range, so BigQuery scans only matching partitions rather than the whole table. This directly cuts bytes processed for time-range queries, satisfying the cost-optimisation requirement in the stem.

Why this answer

Option B is correct because partitioning the table by date lets BigQuery prune partitions and scan only the date ranges a query touches, which directly reduces bytes processed and therefore cost for a dataset that is frequently queried by time range. Option C is correct because committed use reservations (capacity commitments) provide dedicated, predictable slot capacity at a discounted rate for steady baseline workloads, satisfying the requirement that predictable workloads have dedicated slots. Option A is not appropriate as a primary cost strategy here because query caching only helps when identical queries are repeated and results are still cached, which does not address time-range pruning or dedicated capacity.

Option D is wrong because a materialized view over the entire table would be expensive to maintain and does not target the time-range access pattern. Option E is wrong because autoscaling slots add flexible, on-demand capacity rather than the dedicated, predictable capacity the scenario requires.

Exam trap

The trap is picking autoscaling slots for 'predictable workloads' — autoscaling is for variable demand, while committed use reservations are the correct answer for dedicated, predictable baseline capacity.

49
Multi-Selectmedium

Your company uses Cloud Composer to orchestrate a data pipeline that includes Dataproc Spark jobs and BigQuery load operations. You need to pass the output file path from the Spark job to the next BigQuery task in the DAG. Which two mechanisms can you use to share data between tasks? (Choose TWO.)

Select 2 answers
A.Store the output path as a Cloud Composer variable.
B.Publish the output path to a Pub/Sub topic and subscribe in the next task.
C.Write the output path to a Cloud Storage object and read it in the next task.
D.Use BigQuery as an intermediary to store the output path.
E.Use Airflow XComs to push the output path from the Spark task and pull it in the BigQuery task.
AnswersC, E

Writing the output path to a Cloud Storage object and reading it in the next task decouples the Dataproc Spark job from the BigQuery load, persisting the value outside the DAG run. This satisfies the requirement to pass the file path reliably between tasks.

Why this answer

Option E is correct because Airflow XComs are the native, built-in mechanism for passing small pieces of data such as a file path between tasks in a DAG: the Spark task calls xcom_push (or returns a value) and the BigQuery task uses xcom_pull with the task_id/key to retrieve it. Option C is correct because writing the output path to a Cloud Storage object and reading it in the next task is a durable, decoupled way to share data between tasks, especially when the value is larger or must persist beyond the DAG run; Cloud Composer tasks can use the GCS operators/sensors or the google-cloud-storage client to write and read the object. Option A is not appropriate because Cloud Composer (Airflow) Variables are global key-value configuration stores, not per-run task-to-task data channels, so using them for a run-specific output path risks race conditions and stale values.

Option B is not appropriate because Pub/Sub is an asynchronous messaging service, not a synchronous task-to-task handoff mechanism, and it adds unnecessary complexity and latency for passing a single path within one DAG run. Option D is not appropriate because using BigQuery as an intermediary to store a transient output path is heavyweight and semantically wrong; BigQuery is for analytical data, not for orchestrating task metadata.

Exam trap

PDE often tests the misconception that global variables or external services like Pub/Sub are appropriate for task-to-task communication, when in fact Airflow XComs and Cloud Storage are the canonical mechanisms.

50
MCQmedium

Your organization runs a Cloud Composer (Apache Airflow) environment that executes a DAG every night to move data from Cloud Storage into BigQuery. The DAG has been succeeding for months, but last week the nightly load silently produced a BigQuery table with zero rows while the Airflow task still reported success. You need to add a safeguard that fails the DAG task whenever the loaded row count is zero before downstream tasks run. What should you do?

A.Increase the retries parameter on the load task so transient issues that produce zero rows are automatically retried.
B.Add a BigQueryCheckOperator task immediately after the load task that asserts the target table row count is greater than zero.
C.Enable Cloud Logging export for the Composer environment and create a log-based alert that matches error text from the load operator.
D.Configure an Airflow SLA on the load task so that a missed deadline raises an alert about the empty table.
AnswerB

The BigQueryCheckOperator runs a SQL statement and fails the task when the returned result does not match the expected condition. Placing it right after the load task means a zero-row result raises an Airflow task failure, stopping downstream tasks and surfacing the problem through the normal DAG alerting path instead of silently completing.

Why this answer

The failure mode is a task that reports success while producing an empty result set, so the safeguard must inspect the data itself and raise a task-level failure. A BigQueryCheckOperator performs exactly that assertion and integrates with Airflow's dependency graph, which prevents downstream tasks from consuming invalid output and triggers configured alerting.

Exam trap

The trap here is assuming that any task marked successful in Airflow implies the underlying data operation was semantically correct, when success only reflects that the operator did not raise an exception.

51
MCQhard

A data engineer needs to alert when Pub/Sub subscription has messages older than 1 hour. Which Cloud Monitoring metric and filter should they use?

A.Metric: topic/send_message_operation_count; filter: topic_id
B.Metric: subscription/ack_message_count; filter: subscription_id
C.Metric: subscription/num_undelivered_messages; filter: subscription_id
D.Metric: subscription/oldest_unacked_message_age; filter: subscription_id
AnswerD

The `subscription/oldest_unacked_message_age` metric reports the age of the oldest unacknowledged message, directly satisfying the one-hour staleness threshold. Filtering by `subscription_id` scopes the alert to the specific Pub/Sub subscription, so Cloud Monitoring evaluates only that subscription's backlog age rather than aggregating across the project.

Why this answer

The metric subscription/oldest_unacked_message_age directly reports the age of the oldest unacknowledged message in a subscription, which is exactly what's needed to alert when messages are older than 1 hour. The filter subscription_id scopes the metric to the specific subscription. This metric is designed for monitoring message backlog age and is the correct choice for this requirement.

Exam trap

PDE often tests the distinction between message count and message age metrics, tempting candidates to choose num_undelivered_messages when the requirement is about age.

How to eliminate wrong answers

Option A is wrong because topic/send_message_operation_count measures the number of publish operations to a topic, not message age in a subscription. Option B is wrong because subscription/ack_message_count counts acknowledged messages, not the age of unacknowledged ones. Option C is wrong because subscription/num_undelivered_messages gives the count of undelivered messages, not their age, so it cannot detect messages older than 1 hour.

52
MCQmedium

You are monitoring a Dataproc cluster and notice that the cluster utilisation is high, but jobs are running slowly. The cluster uses preemptible workers for cost savings. What is the most likely cause of the performance degradation?

A.The primary workers are using standard disks instead of SSDs.
B.The preemptible workers are being preempted frequently, causing task retries and slowdowns.
C.The cluster is under-provisioned; increase the number of preemptible workers.
D.The cluster is using an older image version; upgrade to the latest.
AnswerB

Frequent preemption of preemptible VMs forces Dataproc to rerun interrupted tasks on remaining workers, inflating job duration despite high CPU utilisation. Since the stem specifies preemptible workers used for cost savings, this directly explains the slowdown: capacity vanishes mid-job, triggering retries and stragglers rather than a genuine resource shortage.

Why this answer

Preemptible (spot) workers in Dataproc can be reclaimed by Compute Engine at any time, and when they are preempted mid-task, the task must be retried on remaining workers. Frequent preemption causes repeated task retries, stragglers, and overall job slowdown even though the cluster appears highly utilized. This is the classic symptom of relying heavily on preemptible workers.

Exam trap

PDE often tests the misconception that high utilization means the cluster is healthy, when in fact preemptible worker churn causes retries that inflate utilization while slowing jobs.

How to eliminate wrong answers

Option A is wrong because standard disks on primary workers would slow I/O but would not produce the intermittent, retry-driven slowdown pattern associated with preemption; also, primary workers are not preempted. Option C is wrong because adding more preemptible workers does not fix the root cause — more preemptible capacity can actually increase the preemption rate and retry churn. Option D is wrong because an older image version is a static configuration issue that would affect all jobs consistently, not cause the fluctuating slowdowns tied to preemption.

53
MCQeasy

A data engineering team maintains several Cloud Composer DAGs that load files from Cloud Storage into BigQuery. They want the DAG to start only after a new file has landed in the source bucket, rather than polling on a fixed schedule and often finding nothing. Which Airflow construct should they use to trigger the DAG based on the arrival of the object?

A.A GCSObjectExistenceSensor with a poke interval, followed by the load task.
B.A TimeDeltaSensor set to the expected upload window before the load task.
C.A BigQueryInsertJobOperator configured to load the file directly from the bucket path.
D.An ExternalTaskSensor pointed at a separate file-upload DAG.
AnswerA

The GCSObjectExistenceSensor waits until the specified object exists in Cloud Storage before allowing the DAG to proceed. This turns a time-driven schedule into an event-aware one: the load task runs only when the file is actually present, eliminating wasted polling runs that find nothing and reducing BigQuery load attempts against empty inputs.

Why this answer

The requirement is to begin processing only when a specific object appears, which is precisely what a Cloud Storage sensor provides. Sensors defer downstream tasks until their condition is met, so pairing a GCSObjectExistenceSensor with the load task converts the DAG from a blind schedule into a data-arrival trigger while keeping the existing load logic intact.

Exam trap

The trap here is confusing a time-based wait with an event-based wait, since both keep a task pending, but only a sensor that inspects Cloud Storage can confirm the object actually exists.

54
MCQeasy

A data engineer needs to run a recurring nightly extract-transform-load job that pulls data from a REST API, applies Python transformations, and writes the output to a Cloud Storage bucket. The team wants a fully managed, serverless scheduler that can retry failed runs and send notifications, and they do not want to maintain any cluster or VM. Which Google Cloud service should they use to define and run this job?

A.A Compute Engine instance running a cron job that executes a Python script and writes to Cloud Storage.
B.Cloud Composer with a single-node environment and a KubernetesPodOperator that runs the Python transformation.
C.Cloud Scheduler triggering a Cloud Function that performs the API call, transformation, and Cloud Storage write.
D.Dataproc Serverless with a PySpark batch that calls the REST API and writes results to Cloud Storage.
AnswerC

Cloud Scheduler provides a fully managed cron service that can trigger an HTTP or Pub/Sub target on a schedule, and Cloud Functions runs the Python code serverlessly with automatic retries and logging. This combination meets the requirement of no cluster or VM to maintain, supports retries on failure, and can publish to a Pub/Sub topic for notifications. It is the lightest managed option for a recurring single-step ETL job.

Why this answer

Cloud Scheduler plus Cloud Functions delivers a fully managed, serverless combination for a recurring ETL task. Cloud Scheduler handles the cron cadence and can target a function over HTTP or Pub/Sub, while Cloud Functions executes the Python code with automatic retries, integrated logging, and the ability to publish notifications. No cluster or VM is required, satisfying the operational constraints.

Exam trap

The trap here is treating Dataproc Serverless as a complete scheduling solution, when it only runs the compute and still needs an external trigger for a nightly cadence.

55
MCQmedium

You manage a Dataflow streaming pipeline that reads from Pub/Sub and writes to BigQuery. The pipeline must be updated to add a new transformation that enriches each message with data from Cloud SQL. You need to minimize downtime and ensure the pipeline continues processing without data loss. What should you do?

A.Use Cloud Scheduler to trigger a Cloud Function that updates the pipeline code in place.
B.Create a new Pub/Sub subscription, update the pipeline to read from it, and deploy a new pipeline.
C.Stop the existing pipeline, update the code, and start a new pipeline from the same subscription.
D.Use the Dataflow update feature to replace the pipeline with the new code while preserving the pipeline state.
AnswerD

The Dataflow update feature allows you to replace the running pipeline with a new version while maintaining the pipeline's state, such as Pub/Sub checkpoints and BigQuery write buffers. This minimizes downtime and ensures no data loss. It is the recommended approach for updating streaming pipelines in production when the update is compatible with the existing pipeline structure.

Why this answer

The Dataflow update feature is designed to replace a running pipeline with a new version while preserving state, enabling seamless updates with minimal downtime and no data loss. Stopping and restarting, creating new subscriptions, or using external triggers do not maintain pipeline state and risk data duplication or loss.

Exam trap

The trap here is assuming that stopping and restarting a pipeline from the same subscription is safe, but it can lead to duplicate or lost messages without proper checkpointing.

56
MCQmedium

A company runs a critical batch pipeline using Cloud Dataflow. The pipeline processes financial transactions and runs every hour. Recently, some runs have failed due to transient errors (e.g., network timeouts). The engineer wants to automatically retry failed runs without manual intervention. The pipeline is launched from a Cloud Composer DAG using DataflowPythonOperator. What is the BEST way to handle retries?

A.Add a DataflowJobStatusSensor in the DAG that waits for job completion and retries if failed.
B.Set the 'retries' parameter in the DAG's default_args to a positive integer.
C.Configure the Dataflow pipeline to automatically retry on failure using the --numberOfWorkerHarnessThreads option.
D.Use a Cloud Function triggered by Cloud Scheduler to re-launch the pipeline if the Dataflow job fails.
AnswerB

Setting retries in the DAG's default_args makes Cloud Composer automatically re-run the DataflowPythonOperator task after transient failures such as network timeouts, without manual intervention. Airflow's task-level retry mechanism satisfies the requirement for unattended recovery of hourly runs.

Why this answer

Cloud Composer (Apache Airflow) natively supports task-level retries via the 'retries' parameter in default_args. When a DataflowPythonOperator fails due to a transient error, Airflow automatically re-executes the task up to the specified number of retries, without requiring custom sensors or external triggers. This is the simplest and most reliable mechanism for handling transient failures in a DAG-driven pipeline.

Exam trap

The trap here is that candidates confuse Dataflow-level retry options (like --maxRetryAttempts) with Airflow task-level retries, or assume that a sensor or external trigger is required to detect and retry failures, when in fact Airflow's native retry parameter is the simplest and most appropriate solution for transient errors in a DAG-managed pipeline.

How to eliminate wrong answers

Option A is wrong because a DataflowJobStatusSensor only monitors job status and does not automatically retry the pipeline; it would require additional branching logic to relaunch the job, adding unnecessary complexity. Option C is wrong because --numberOfWorkerHarnessThreads controls parallelism within the Dataflow worker, not retry behavior on pipeline failure; retries are configured via --maxRetryAttempts or similar Dataflow pipeline options, not this flag. Option D is wrong because using a Cloud Function and Cloud Scheduler introduces an external dependency and latency, whereas Airflow's built-in retry mechanism is more direct and integrated with the DAG's execution context.

57
MCQeasy

An organization uses BigQuery on-demand pricing. To control costs, they want to estimate the bytes processed by a query before running it. Which command or method should they use?

A.Use the bq query --dry_run command
B.Use bq ls to list table sizes
C.Use BigQuery reservations to get cost estimate
D.Use INFORMATION_SCHEMA.JOBS_BY_PROJECT to view past costs
AnswerA

The `bq query --dry_run` flag validates a query and returns the bytes it would process without executing it, so no on-demand charges are incurred. This directly satisfies the requirement to estimate bytes processed beforehand, letting the organisation predict cost before committing to the query.

Why this answer

The bq query --dry_run flag validates a query and returns the estimated bytes that would be processed without actually executing it or incurring charges. This is the standard method to estimate on-demand BigQuery costs before running a query, since on-demand pricing is based on bytes processed.

Exam trap

PDE often tests the difference between pre-execution estimation and post-hoc cost inspection, so the trap is choosing INFORMATION_SCHEMA.JOBS_BY_PROJECT (which shows past bytes billed) instead of --dry_run for estimating before running.

How to eliminate wrong answers

Option B is wrong because bq ls only lists datasets, tables, or jobs—it does not estimate query bytes processed. Option C is wrong because reservations are for capacity-based (flat-rate) pricing and do not provide a per-query byte estimate. Option D is wrong because INFORMATION_SCHEMA.JOBS_BY_PROJECT shows historical job metadata and bytes billed after the fact, not a pre-execution estimate.

58
MCQmedium

You have a BigQuery table that is used by multiple teams. To save costs, you want to provide a consistent view of the data as of a specific point in time without creating full copies. Which BigQuery feature should you use?

A.Authorized views
B.Materialized views
C.Table snapshot
D.Table clone
AnswerC

A table snapshot captures the table's data at a specific point in time using a copy-on-write reference, so it is queryable and cheap to create without duplicating storage. This gives multiple teams a consistent historical view while avoiding the cost of full table copies.

Why this answer

A BigQuery table snapshot creates a lightweight, read-only copy of a table at a specific point in time without duplicating the underlying data, and it is the correct feature for providing a consistent point-in-time view cheaply. Snapshots preserve the table state at creation and can be queried or restored, making them ideal for this requirement.

Exam trap

PDE often tests the difference between snapshots (read-only, point-in-time, copy-on-write) and clones (writable, diverging), so the trap is choosing table clone or materialized view when the requirement is a frozen, consistent point-in-time view.

How to eliminate wrong answers

Option A is wrong because authorized views control access to data, not point-in-time consistency, and do not create a snapshot of the data state. Option B is wrong because materialized views precompute and cache query results for acceleration, not for preserving a point-in-time copy, and they refresh as the base table changes. Option D is wrong because a table clone creates a writable copy that shares storage initially but diverges on writes and does not freeze a point-in-time state.

59
Multi-Selecteasy

You need to deploy a reusable Dataflow pipeline that can be executed with different parameters from Cloud Composer. Which TWO components should you use? (Choose 2)

Select 2 answers
A.Direct runner
B.Dataflow Flex Template
C.Cloud Composer with DataflowStartFlexTemplateOperator
D.Dataflow Classic Template
E.Cloud Scheduler
AnswersB, C

A Flex Template packages the pipeline as a Docker image with a metadata parameter spec, so the same artefact runs repeatedly with different runtime parameters. This satisfies the reusability constraint, unlike classic templates, which require rebuilding or redeploying the pipeline for parameter changes.

Why this answer

Option B (Dataflow Flex Template) is correct because Flex Templates package the pipeline as a Docker image plus a template spec file, allowing the same pipeline to be reused and launched with different runtime parameters (such as input/output locations) via the templates launch API. Option C (Cloud Composer with DataflowStartFlexTemplateOperator) is correct because this operator is purpose-built to submit a Flex Template job from an Airflow/Composer DAG, passing the required parameters and letting Composer orchestrate the pipeline execution. Option A (Direct runner) is incorrect because it runs the pipeline locally for testing rather than deploying a reusable job on Dataflow.

Option D (Dataflow Classic Template) is not the best fit here since Classic Templates are less flexible for parameterization and dependency packaging compared with Flex Templates. Option E (Cloud Scheduler) is incorrect because it only triggers jobs on a schedule and does not deploy or parameterize a Dataflow pipeline.

Exam trap

PDE often tests the distinction between Flex Templates and Classic Templates, as candidates may choose Classic Templates for reusability but overlook the need for runtime parameterization and the specific operator for Cloud Composer integration.

60
MCQhard

You manage a Cloud Composer environment that runs a critical DAG every hour. The DAG includes a task that calls a Cloud Function to process data. Recently, the Cloud Function started taking longer than expected, causing the DAG to exceed its SLA. You need to detect this delay and automatically retry the task if it fails due to timeout, while minimizing changes to the DAG. What should you do?

A.Use a Cloud Monitoring alert on the Cloud Function's execution time and manually trigger the DAG upon alert.
B.Increase the Cloud Function's timeout setting to the maximum allowed and rely on the DAG's default retry behavior.
C.Set the task's execution_timeout parameter to a value slightly above the expected runtime and configure retries with a delay.
D.Add a Python operator that polls the Cloud Function's status and raises an exception if it exceeds a threshold, then set retries on that operator.
AnswerC

The execution_timeout parameter defines the maximum time a task can run before it is killed and marked as failed. Setting it appropriately allows the task to be retried if it exceeds the timeout. Configuring retries ensures automatic retry on failure. This approach requires minimal changes to the DAG and directly addresses the timeout issue, making it the most efficient solution.

Why this answer

Using the task's execution_timeout parameter allows Airflow to enforce a time limit on the task. If the task exceeds this limit, it fails and can be automatically retried based on the retries and retry_delay settings. This requires only a small change to the DAG's task definition and leverages native Airflow features, effectively addressing both detection of delays and automatic retries.

Exam trap

The trap here is focusing on the Cloud Function's timeout instead of the Airflow task's execution_timeout, which controls retries at the orchestration level.

61
MCQeasy

A data engineer needs to inspect a BigQuery table for sensitive data such as credit card numbers and email addresses before sharing it with a third party. The engineer also wants to de-identify the data by masking the sensitive columns. Which Google Cloud service should be used?

A.Dataplex
B.BigQuery column-level security
C.Data Catalog
D.Cloud DLP
AnswerD

Cloud DLP inspects BigQuery data for infoTypes such as credit card numbers and email addresses, then applies de-identification transforms like masking or tokenisation. It satisfies both requirements in one service, discovery and masking, before the table is shared with the third party.

Why this answer

Cloud DLP (Data Loss Prevention) is the correct service because it is specifically designed to inspect, classify, and de-identify sensitive data such as credit card numbers and email addresses. It provides built-in infoType detectors for over 150 types of sensitive data and supports masking, tokenization, and other de-identification techniques. The engineer can use Cloud DLP to scan BigQuery tables and then apply transformations to mask the sensitive columns before sharing.

Exam trap

The trap here is that candidates confuse BigQuery column-level security (access control) with data masking, but BigQuery column-level security only hides data from unauthorized users and does not inspect or transform the data itself, whereas Cloud DLP performs actual de-identification.

How to eliminate wrong answers

Option A is wrong because Dataplex is a data fabric service for managing, governing, and cataloging data across lakes, warehouses, and marts, but it does not natively inspect or de-identify sensitive data; it relies on integration with Cloud DLP for such tasks. Option B is wrong because BigQuery column-level security (using policy tags) only controls access at the column level by granting or denying read permissions, but it does not inspect content for sensitive data or perform masking/de-identification. Option C is wrong because Data Catalog is a metadata management service for discovering and tagging data assets, but it cannot scan for sensitive data patterns or apply de-identification transformations.

62
MCQmedium

A company uses Cloud Pub/Sub for a real-time data pipeline. The subscription has a backlog of millions of messages that are not being processed quickly enough. In Cloud Monitoring, you observe that the 'subscription/num_undelivered_messages' metric is high and growing, while 'subscription/oldest_unacked_message_age' is also increasing. Which action is MOST likely to reduce the backlog?

A.Delete the subscription and recreate it with a larger message retention duration.
B.Reduce the acknowledgment deadline to force faster processing.
C.Change the subscription type from push to pull.
D.Increase the number of subscribers or the throughput capacity of the existing subscribers.
AnswerD

Scaling subscriber count or per-subscriber throughput directly raises the rate at which messages are pulled and acknowledged, so the undelivered-message count and oldest-unacked age both fall. The backlog stems from consumption capacity lagging behind publish rate, not from retention or delivery configuration, so adding consumer capacity addresses the actual constraint.

Why this answer

The backlog indicates that subscribers cannot keep up with the message flow. Increasing the number of subscribers or scaling their throughput capacity directly addresses the processing bottleneck, allowing messages to be pulled and acknowledged faster. Cloud Pub/Sub scales horizontally, so adding more pull subscribers or increasing the resources of existing ones (e.g., more worker threads, higher CPU/memory) reduces the backlog.

Exam trap

A common mistake is to think that reducing the acknowledgment deadline or changing subscription type will speed up processing, when in reality these actions can increase redeliveries or do not address the root cause of insufficient subscriber capacity.

How to eliminate wrong answers

Option A is wrong because deleting and recreating the subscription with a larger message retention duration does not increase processing speed; it only keeps messages longer, which does not reduce the existing backlog. Option B is wrong because reducing the acknowledgment deadline forces subscribers to acknowledge messages faster, but if they cannot process them in time, it leads to more redeliveries and can worsen the backlog. Option C is wrong because changing from push to pull does not inherently increase throughput; both modes can be scaled, and the bottleneck is subscriber capacity, not the delivery mechanism.

63
MCQmedium

A data platform team wants to grant a service account the ability to run BigQuery jobs and read data in a specific dataset, while ensuring it cannot create or delete datasets. Which IAM approach satisfies this with least privilege?

A.Grant the BigQuery Admin role on the dataset.
B.Grant the BigQuery Data Editor role at the project level.
C.Grant bigquery.jobs.create at the project level and BigQuery Data Viewer on the specific dataset.
D.Grant the service account the BigQuery Data Viewer role at the project level.
AnswerC

Job creation permission is granted at the project level because jobs are project-scoped resources, while data read access is granted on the dataset to limit exposure. This combination lets the service account run jobs that read the target dataset without granting rights over other datasets. It excludes dataset creation and deletion permissions, satisfying least privilege.

Why this answer

BigQuery separates the permission to run jobs from the permission to read data. Job creation is a project-level capability, so the service account needs bigquery.jobs.create at the project scope. Reading data can be scoped to the specific dataset with BigQuery Data Viewer, which confines access to that dataset and excludes dataset creation or deletion.

Combining these two grants achieves the read-and-run requirement with least privilege.

Exam trap

The trap here is assuming a single predefined role covers both running jobs and reading a dataset, when job creation and data access are granted at different scopes.

64
MCQhard

You are optimizing a BigQuery query that scans 1 TB of data every day. The query joins a large fact table (partitioned by date) with a small dimension table. You notice that the query always scans the entire fact table, even though you only need the last 7 days of data. Which optimization will MOST reduce the bytes scanned?

A.Create a materialized view that pre-aggregates the data by day.
B.Add a WHERE clause that filters on the date column used for partitioning.
C.Cluster the fact table on the join key used in the query.
D.Change the table to use time-unit column partitioning with a 1-day partition interval.
AnswerB

Partition pruning only activates when the query filters on the partitioning column. Adding that WHERE predicate lets BigQuery eliminate all but the last seven date partitions, cutting bytes scanned from 1 TB to roughly 7 days' worth.

Why this answer

BigQuery partitioned tables allow partition pruning, where the query engine skips partitions that do not match the filter condition. Adding a WHERE clause on the partitioning column (date) restricts the scan to only the last 7 days' partitions, drastically reducing bytes scanned. This is the most direct and effective optimization for the described scenario.

Exam trap

PDE often tests whether candidates know that partition pruning requires a filter on the partitioning column; candidates may choose clustering or materialized views, but the most direct reduction in bytes scanned comes from adding the WHERE clause on the partition column.

How to eliminate wrong answers

Option A is wrong because a materialized view pre-aggregates data but does not reduce the bytes scanned for the underlying fact table unless the query is rewritten to use the view; it also adds storage and refresh costs. Option C is wrong because clustering on the join key improves filter and join performance but does not eliminate scanning of all partitions if the date filter is absent; clustering helps with selective filters on clustered columns, but partition pruning is more effective for date ranges. Option D is wrong because changing to time-unit column partitioning with a 1-day interval is already likely the case (partitioned by date), and it does not by itself reduce bytes scanned without a filter; the partitioning method is not the issue, the missing filter is.

65
MCQmedium

You are running a streaming pipeline with Dataflow that reads from Pub/Sub and writes to BigQuery. You notice that the system lag metric is increasing over time, indicating that messages are taking longer to process. What is the most likely cause and how should you address it?

A.The source Pub/Sub topic has insufficient throughput; increase the number of partitions.
B.The Dataflow workers are CPU-bound; increase the number of workers or adjust autoscaling settings.
C.The BigQuery destination table has too many columns; reduce the number of columns.
D.The pipeline uses a batch transform that should be replaced with a streaming transform.
AnswerB

High system lag suggests worker resources are insufficient; adding workers reduces lag.

Why this answer

System lag in Dataflow measures the age of the oldest unprocessed message, and a steadily increasing value indicates the pipeline cannot keep up with input. The most common cause is CPU-bound workers, so scaling out workers or tuning autoscaling (e.g., max workers, worker type) restores throughput. This directly addresses the bottleneck rather than a downstream symptom.

Exam trap

The trap is picking a source-side fix (partitions) that applies to Kafka, not Pub/Sub, or blaming the sink; the exam expects you to recognize that increasing system lag in Dataflow is typically a worker capacity issue solved by autoscaling.

How to eliminate wrong answers

Option A is wrong because Pub/Sub does not use partitions; it uses subscriptions and scales via the number of messages and ack deadlines, so 'increase partitions' is a Kafka concept misapplied here. Option C is wrong because BigQuery column count does not cause processing lag; BigQuery handles wide tables efficiently and the bottleneck would be elsewhere. Option D is wrong because the pipeline is already streaming (Pub/Sub to BigQuery), so there is no batch transform to replace.

66
MCQmedium

A data platform team uses Cloud Composer 2 to orchestrate a DAG that runs a Dataproc Serverless batch job producing a partitioned BigQuery table. The DAG passes the output location to a downstream task that runs a dbt model. The team wants failures in the dbt task to automatically trigger a retry of only that task, and they want the DAG to expose the Dataproc job ID in the Airflow UI for troubleshooting. Which approach BEST satisfies both requirements?

A.Use the DataprocSubmitJobOperator and configure the DAG's default_args with retries, which applies the retry policy to all tasks including the dbt task.
B.Use a single PythonOperator that submits the Dataproc job via the API, waits for completion, and then invokes dbt, with retries set on the PythonOperator.
C.Use the DataprocSubmitJobOperator and pass the returned job reference to XCom, then set retries on the dbt task and read the XCom value in its template fields.
D.Use the DataprocSubmitJobOperator with retries set on it, and add a separate downstream PythonOperator that runs dbt without its own retry configuration.
AnswerC

DataprocSubmitJobOperator returns the job reference, which Airflow automatically pushes to XCom, making the job ID visible and available downstream. Setting retries on the dbt task ensures only that task is retried on failure, and using XCom in template fields lets the dbt task consume the job ID or output path. This satisfies both the observability and granular retry requirements.

Why this answer

DataprocSubmitJobOperator returns a job reference that Airflow pushes to XCom, which makes the job ID available for downstream tasks and visible in the Airflow UI. Configuring retries specifically on the dbt task ensures that a transient failure there triggers only that task to rerun, leaving the already-successful Dataproc submission intact. This combination precisely meets both the observability and granular retry goals.

Exam trap

The trap here is setting retries on the upstream Dataproc operator or in default_args, when the requirement is to retry only the failing dbt task.

67
MCQmedium

Your team uses Cloud Dataproc for Spark ML training jobs. You want to reduce costs for non-critical, fault-tolerant training jobs. Which Dataproc feature should you use for worker nodes?

A.Use preemptible instances for worker nodes.
B.Use custom machine types with more memory.
C.Use SSDs instead of HDDs for persistent disks.
D.Use committed use discounts for 1-year or 3-year terms.
AnswerA

Preemptible instances suit fault-tolerant Spark training because Dataproc replaces them automatically when Compute Engine reclaims capacity, and they cost substantially less than standard VMs. The stem's non-critical, fault-tolerant constraint is satisfied: interrupted workers are re-added without failing the job, unlike sole-primary or non-retryable workloads.

Why this answer

Preemptible instances are short-lived, lower-cost VMs that Cloud Dataproc can use for worker nodes. Because the training jobs are non-critical and fault-tolerant (e.g., they can handle node failures via Spark's built-in resilience), preemptible instances significantly reduce costs while still completing the workload. This directly addresses the requirement to reduce costs for fault-tolerant jobs.

Exam trap

A common pitfall is confusing committed use discounts (which require a 1- or 3-year commitment) with preemptible VMs (which are interruptible but cost-effective). For non-critical, fault-tolerant workloads, preemptible VMs are the appropriate cost-saving feature, not committed use discounts.

How to eliminate wrong answers

Option B is wrong because custom machine types with more memory increase cost per node, which contradicts the goal of reducing costs. Option C is wrong because SSDs are more expensive than HDDs, and while they improve I/O performance, the question focuses on cost reduction, not performance. Option D is wrong because committed use discounts require a 1-year or 3-year commitment and are typically applied to all instances in a project, not specifically to worker nodes in a Dataproc cluster; they also do not leverage the fault-tolerant nature of the jobs to achieve the lowest possible cost.

68
MCQeasy

A data engineer needs to orchestrate a complex data pipeline that involves multiple steps including data extraction from Cloud Storage, transformation using Dataflow, and loading into BigQuery. The pipeline has dependencies between tasks and requires monitoring and retries. Which Google Cloud service should be used for orchestration?

A.Workflows
B.Cloud Scheduler
C.Cloud Composer
D.Cloud Tasks
AnswerC

Cloud Composer provides managed Apache Airflow, whose directed acyclic graphs express task dependencies, scheduling, retries and monitoring across extraction, Dataflow transformation and BigQuery loading. This satisfies the stem's requirement for orchestration of a multi-step pipeline with dependencies and retry handling.

Why this answer

Cloud Composer is a fully managed workflow orchestration service built on Apache Airflow. It is designed to orchestrate complex data pipelines with dependencies, scheduling, monitoring, and retries, making it the correct choice for orchestrating extraction, transformation, and loading tasks across Google Cloud services.

Exam trap

PDE often tests the distinction between orchestration services; candidates may choose Workflows because it is also an orchestration tool, but Cloud Composer is specifically designed for data pipelines with Airflow.

How to eliminate wrong answers

Option A is wrong because Workflows is a serverless orchestration service for HTTP-based APIs and services, but it lacks the rich scheduling and dependency management of Airflow for data pipelines. Option B is wrong because Cloud Scheduler is a cron job service that triggers jobs but does not orchestrate multi-step pipelines with dependencies. Option D is wrong because Cloud Tasks is a task queue service for asynchronous task execution, not a pipeline orchestration tool.

69
MCQhard

A data engineer manages a Cloud Composer 2 environment. A DAG that downloads a large reference dataset each night occasionally exceeds the default task timeout because the source API is slow. The engineer wants the task to fail fast and be retried automatically rather than hanging for hours, and wants failed runs to be visible for alerting. Which configuration should be applied to the task?

A.Configure the worker's celery_worker_autoscale setting to add more workers so the slow API call finishes sooner.
B.Set depends_on_past to true and add a retry_exponential_backoff so the task waits for the prior run before starting.
C.Increase the DAG's schedule interval and add a max_active_runs limit of one so overlapping runs cannot occur.
D.Set execution_timeout on the task to a bounded value and set retries with a retry_delay so failed attempts are re-queued and surfaced in task state.
AnswerD

execution_timeout caps how long a task may run before Airflow marks it failed, and the retries plus retry_delay settings cause the scheduler to re-queue it automatically. The failure state is recorded and available for alerting, which matches the requirement to fail fast, retry, and remain visible.

Why this answer

execution_timeout is the task-level setting that bounds wall-clock runtime and triggers a failure, while retries and retry_delay make the scheduler automatically re-queue the task. Together they produce fast failure, automatic retry, and a recorded state that alerting can observe, which is precisely what the scenario requires for a slow external API.

Exam trap

The trap here is confusing concurrency controls such as max_active_runs or worker autoscaling with a per-task execution timeout.

70
MCQhard

A company wants to use Cloud DLP to inspect data in BigQuery for sensitive information and de-identify it by masking credit card numbers. They want to perform this on a schedule. Which approach should they take?

A.Use Dataplex data quality rules with a custom SQL regex
B.Use Cloud Data Loss Prevention API with Cloud Composer
C.Use BigQuery column-level security with classification
D.Use Cloud DLP inspect and de-identify jobs triggered by Cloud Scheduler
AnswerD

Cloud Scheduler triggers Cloud DLP inspect and de-identify jobs on a defined cadence, satisfying the scheduled requirement. De-identification uses the masking transformation to obscure credit card numbers while preserving data format. This combination inspects BigQuery data and applies masking without manual intervention, matching both the recurring schedule and de-identification constraints in the stem.

Why this answer

Cloud DLP can inspect BigQuery tables and de-identify using transforms like masking. Scheduling can be done via Cloud Scheduler.

71
Multi-Selectmedium

You are building a data pipeline that ingests data from on-premises into Cloud Storage, then processes it with Dataproc, and finally loads into BigQuery. You need to schedule the pipeline to run daily. The pipeline must handle occasional failures gracefully. Which THREE Google Cloud services should you use together to achieve this? (Choose 3)

Select 3 answers
A.Cloud Storage
B.Dataproc
C.Cloud Composer
D.Dataflow
E.Pub/Sub
AnswersA, B, C

Cloud Storage serves as the ingestion landing zone for on-premises data before Dataproc processing, satisfying the pipeline's first stage. It provides durable, highly available object storage that decouples ingestion from compute, letting Dataproc read source files and write results without data loss during the daily scheduled runs.

Why this answer

Cloud Storage (A) is correct because the scenario explicitly requires ingesting on-premises data into Cloud Storage as the landing/staging layer before processing. Dataproc (B) is correct because the pipeline is specified to process the data with Dataproc, Google Cloud's managed Hadoop/Spark service. Cloud Composer (C) is correct because it is the managed Apache Airflow service that provides daily scheduling, dependency orchestration, and retry/alerting logic so occasional failures are handled gracefully.

Dataflow (D) is not needed since the processing engine in this pipeline is Dataproc, not a Beam-based Dataflow job. Pub/Sub (E) is not required because the scenario describes batch ingestion into Cloud Storage rather than real-time streaming messaging.

Exam trap

PDE often tests whether candidates confuse orchestration (Composer/Airflow) with processing (Dataflow/Dataproc) — picking Dataflow because it 'processes data' ignores that the question already specifies Dataproc and asks for scheduling.

72
Multi-Selecthard

A media company runs a Cloud Composer environment whose DAGs trigger Dataflow batch jobs and BigQuery loads. The operations team reports that Composer costs are rising and that DAGs occasionally stall because workers are saturated. You review the environment and find that several tasks are long-running sensors that hold worker slots while waiting on external conditions. Which TWO changes should you make to reduce worker saturation and cost? (Choose two.)

Select 2 answers
A.Move the long-running Dataflow job monitoring out of the DAG by having the DAG submit the job and exit while a separate process tracks completion.
B.Change the DAG schedule from a cron expression to a timedelta interval so runs are spaced further apart.
C.Set the Composer environment's worker count to the minimum and rely on autoscaling to add workers only when the queue is deep.
D.Increase the DAG's max_active_tasks parameter so more tasks can run in parallel on the same workers.
E.Switch the long-waiting sensors to deferrable mode so they release the worker slot while waiting and resume when the condition is met.
AnswersA, E

When a task blocks a worker for the entire duration of a Dataflow job, that slot is unavailable for other work. Submitting the job and returning, then checking status separately or through a deferrable waiter, reduces the time a worker is occupied. This shortens the critical path of worker occupancy and alleviates the saturation the team is seeing.

Why this answer

Worker saturation caused by long-blocking sensors and job-waiting tasks is best relieved by changing how those tasks consume resources rather than by adding concurrency limits or rescheduling. Deferrable sensors release slots during waits, and decoupling long Dataflow job monitoring from the worker keeps the pool available for other tasks, which lowers both stalls and cost.

Exam trap

The trap here is treating worker saturation as a scheduling or concurrency problem, when the real cause is tasks that occupy worker slots for long durations while doing nothing but waiting.

73
MCQhard

A financial services firm runs an Apache Airflow workload on Cloud Composer 2 that ingests market data, runs dbt transformations, and loads curated tables into BigQuery. The DAG currently uses a single PythonOperator that runs a long shell command, and the team wants to make failures easier to diagnose and retries more granular. They also want to avoid rerunning already-successful upstream steps. Which change BEST meets these goals?

A.Move the shell command into a Cloud Function and invoke it from the PythonOperator using an HTTP call, keeping the DAG structure unchanged.
B.Split the monolithic PythonOperator into a chain of task-level operators such as BigQueryInsertJobOperator and DataprocSubmitJobOperator, with retries configured per task and depends_on_past set appropriately.
C.Increase the worker count and worker machine type on the Cloud Composer environment so the single PythonOperator task has more resources and completes faster.
D.Set the DAG's max_active_runs to 1 and enable catchup so that failed runs are automatically retried on the next schedule interval.
AnswerB

Breaking the DAG into discrete task-level operators lets Airflow retry only the failed step and preserves successful upstream results, since each task tracks its own state. Using purpose-built operators such as BigQueryInsertJobOperator gives clearer logs and status per step, and configuring retries at the task level provides granular control. This directly addresses diagnosability and avoids rerunning completed work, which is the goal described.

Why this answer

Decomposing a monolithic task into discrete operators is the standard Airflow pattern for granular retries and clearer observability. Each task maintains its own state, so a failure in the load step can be retried without rerunning ingestion or transformation. Purpose-built operators emit structured logs and status, and task-level retry configuration gives precise control over failure handling.

Exam trap

The trap here is assuming that scaling the Composer environment or enabling catchup improves retry granularity, when those settings affect capacity and scheduling rather than task-level state.

74
MCQmedium

You are designing a Cloud Composer workflow that loads data from Cloud Storage into BigQuery, runs a Dataflow job to transform the data, and then triggers a Dataproc Spark job. After each step, you need to conditionally branch based on success or failure. Which Airflow feature allows you to pass messages between tasks to enable dynamic branching?

A.Sensors
B.XComs
C.TaskFlow API
D.DAG dependencies
AnswerB

XComs let a task push a small value, such as a success flag or branch key, that downstream tasks pull, enabling conditional branching via BranchPythonOperator. This passes messages between tasks without external storage, satisfying the requirement for dynamic branching after each step.

Why this answer

XComs (cross-communications) are Airflow's built-in mechanism for passing small pieces of data between tasks. A task can push a value via xcom_push() or by returning it, and downstream tasks pull it via xcom_pull(), enabling dynamic decisions such as choosing a branch based on a prior task's output. This is exactly what conditional branching in a Composer DAG requires.

Exam trap

PDE often tests the confusion between XComs (data passing) and the TaskFlow API (a coding style that uses XComs) — candidates pick TaskFlow API thinking it is the transport mechanism.

How to eliminate wrong answers

Option A is wrong because Sensors are operators that wait for a condition (a file, a partition, a time) to become true — they do not carry payloads between tasks. Option C is wrong because the TaskFlow API is a decorator-based authoring style (@task) that simplifies DAG code and implicitly uses XComs under the hood; it is not itself the message-passing mechanism. Option D is wrong because DAG dependencies (>>, set_upstream/set_downstream) only define execution order, not data transfer.

75
MCQhard

You are designing the deployment process for a Dataflow streaming pipeline that processes financial transactions. The pipeline must be updated without losing in-flight state, such as open windows and timers, and without downtime. Your team uses the Apache Beam Java SDK and deploys from a CI/CD pipeline. Which update strategy should you use?

A.Submit an update with the --update option and a compatible transform change to replace the pipeline in place.
B.Cancel the pipeline and redeploy with a new job name.
C.Use a snapshot and then start a new job from the snapshot to preserve state.
D.Stop the pipeline with a drain, then start a new pipeline from the same template.
AnswerA

Dataflow supports in-place updates via the update option, which preserves pipeline state including windows and timers if the transform changes are compatible. This allows the pipeline to continue processing without downtime. It is the standard mechanism for evolving streaming pipelines while maintaining exactly-once semantics and state.

Why this answer

Dataflow update replaces a running streaming pipeline in place while preserving state, provided the transform changes are compatible with the update compatibility rules. This avoids downtime and keeps windowed state intact. Draining, snapshot-based restart, and cancel-redeploy all interrupt processing or discard state, making them unsuitable for stateful, zero-downtime updates.

Exam trap

The trap here is assuming a snapshot or drain preserves state for a new job, when only an in-place update maintains state without interrupting processing.

Page 1 of 2 · 83 questions totalNext →

Ready to test yourself?

Try a timed practice session using only Pde Maintaining Automating questions.