managed-airflow-dag-troubleshooting
Provides guidance for troubleshooting Apache Airflow DAGs (failed DAG runs and task instances) in Managed Service for Apache Airflow (MSAA; formerly Cloud Composer). Use when figuring out reasons for DAG run or task instance failures. Don't use when looking for overall recommendations for Managed Ai
- 0
- Installs
- —
- Rating
- —
- Success rate
- 1
- Files scanned
Security scan
Scan passedNo risky patterns were found in the scanned files.
Content sha256 daeef18ac5e94d00… — run codexguild_scan_skills after installing to verify your local copy.
Static analysis is a first line of defense, not a guarantee. Read the source
SKILL.md
Managed Service for Apache Airflow (formerly Cloud Composer) DAG troubleshooting guide
This skill provides instructions for troubleshooting Managed Airflow DAGs (DAG
runs and task instances), utilizing gcloud composer, gcloud logging and
gcloud storage commands to fetch remote logs and code.
General rules
-
Provide suggestions on how to troubleshoot the failed jobs. Provide only the steps that the user can actually take. Ground all troubleshooting advice in direct findings.
-
When troubleshooting a failure, follow the following practices to always provide a deterministic diagnosis:
-
Fetch relevant logs: Always fetch the logs for a task under investigation using
gcloud logging read; check the logs for specific error patterns: Python tracebacks, API error codes (e.g., 400, 403, 404, 500), or Airflow signals (e.g.,AirflowTaskTimeout). -
Fetch task metadata: When troubleshooting a task, fetch the task state and metadata (execution state, try number, timestamps, and execution details) using:
gcloud composer environments run {env_name} \ --location {location} \ tasks states-for-dag-run -- -d {dag_id} -r {run_id}or for an individual task instance:
gcloud composer environments run {env_name} \ --location {location} \ tasks state -- {dag_id} {task_id} {execution_date} -
Retrieve and compare DAG source code: Download the remote DAG source code using
gcloud storage cp gs://{bucket_name}/dags/{dag_file}.py .(find the environment bucket viagcloud composer environments describe {env_name} --location {location} --format="value(config.dagGcsPrefix)"). Compare the parameters in the code (e.g., table IDs, disk sizes, URI paths) against the error messages found in the task logs. -
Explain code mistakes and potential fixes: Explain mistakes in the code (if any are actually visible); suggest potential fixes (if they are very likely to be meaningful); discuss source code availability if needed - if some source code is unavailable (e.g. imported from a file other than the main source code file), mention this (you can mention the package name) - in such a case take into account most likely trigger rules if they are unknown.
-
Check for environment-level errors: Query Cloud Logging with
gcloud logging readto see if there are high-level environment issues or known platform errors correlating with the failure (see Known issues below). You MUST return ALL found issues. -
Identify failing tasks in a DAG run: When troubleshooting a failed DAG run, mention the task that caused a failure (use
tasks states-for-dag-runor Cloud Logging to identify failed tasks). Provide a task instance name. If many tasks failed, mention which task was critical (mandatory for successful DAG run execution - look into task dependencies and trigger rules) and focus on this one. -
Verify service configurations in code: If logs suggest an issue with a specific service (e.g., BigQuery, Dataform, Compute Engine), use the log details to verify the configuration in the DAG source code.
-
Correlate logs with code: E.g., if BigQuery returns a 404, verify the dataset ID or table ID in the DAG source code matches reality.
-
Prioritize known platform issues: Check against Known issues below. If Cloud Logging queries return matching platform error signals, prioritize that diagnosis.
-
-
Summarize with Evidence (Deterministic Response): Your response must be specific. Avoid general advice like 'check your permissions.' or 'check the logs.' Instead, say 'The service account is missing X permission.'
- Problem: State the specific root cause and the exact task instance ID. Identify if it is a code logic error, a configuration mismatch, or an environment timeout.
- Evidence: Mandatory. Provide the verbatim text from the log
(
textPayload) or the specific line of code from the DAG that caused the failure. Do not summarize the evidence; show the data. - Recommendation: Provide an actionable fix. If it is a code error, provide the corrected Python snippet. If it is a resource issue, specify the exact configuration change needed.
-
DAGs Generated by Orchestration Pipelines: Some DAGs may be generated by Orchestration Pipelines. A special requirement related to those DAGs is the need to explain the failure in terms of the logical actions defined in the pipeline YAML.
- Determine if a DAG is generated by Orchestration Pipelines:
Orchestration Pipeline DAGs deployed by dedicated tools have
bundle_name,version_id, andpipeline_nameset in their DAG Run metadata (DagRun.notethat contains JSON metadata). All of them (i.e. Orchestration Pipeline DAGs deployed by dedicated tools and created manually) have anop:orchestration_pipelinetag set (DAG properties, including tags, can be verified in the DAG source code or viagcloud composer environments run {env_name} --location {location} dags list). - Orchestration Pipeline DAGs deployed by dedicated tools have
additionally the following tags (information in those tags should be
consistent with data in DAG Run attributes mentioned above):
- pipeline name - tag
op:pipeline, e.g.op:pipeline:xyzindicates a namexyz - bundle name - tag
op:bundle - version id - tag
op:version
- pipeline name - tag
- Retrieve the resolved pipeline YAML definition from the environment
bucket:
- Determine the YAML file location:
- Retrieve the DAG source code from the environment bucket using
gcloud storage cp gs://{bucket_name}/dags/{dag_file}.py .(orgcloud storage cat gs://{bucket_name}/dags/{dag_file}.py). - Inspect the source code for
generateorgenerate_dagsfunction calls:- Scenario 1:
generatecall found. The first argument is the path to the YAML file - relative to thedagsfolder in environment's bucket. - Scenario 2:
generate_dagscall found.- Extract the first argument - this is the data folder. If
it starts with
/home/airflow/gcs/, remove this prefix to get a path relative to the root of environment's bucket. - Extract
bundle_name,version_id, andpipeline_name(as explained above). - Construct the path:
{data_directory}/{bundle_name}/versions/{version_id}/{pipeline_name}.yml(or.yaml).
- Extract the first argument - this is the data folder. If
it starts with
- Scenario 3: If neither call is found, default to the path:
data/{bundle_name}/versions/{version_id}/{pipeline_name}.yml(or.yaml) in an environment's bucket.
- Scenario 1:
- Download the YAML file using
gcloud storage cp gs://{bucket_name}/{yaml_path} .(orgcloud storage cat gs://{bucket_name}/{yaml_path}).
- Retrieve the DAG source code from the environment bucket using
- Determine the YAML file location:
- Map the failed Airflow task back to the logical action name using task
instance metadata/notes (e.g.
op_action_namein tasknote). - If the failure involves user assets (like Python scripts), check their
path in the action definition. If they are in the environment bucket,
download and read them to debug (
gcloud storage cp gs://{bucket_name}/{asset_path} .). If they are in a custom artifact bucket (see GCS URIs in logs/config), note the limitation that they cannot be read directly but analyze based on available logs.
- Determine if a DAG is generated by Orchestration Pipelines:
Orchestration Pipeline DAGs deployed by dedicated tools have
-
You can assume that environment variables set by default (they can be used in DAG code, but are not visible in custom environment configuration), e.g.
GCS_BUCKET, are correct - users cannot change them. -
"Not found" (404) errors from GCP APIs can be misleading. A "not found" error might be returned when a resource actually exists, but the caller does not have permissions to access or view it. If a resource is expected to exist, suggest verifying proper permissions.
Important constraints & instructions
- Read-Only First: Do NOT attempt to fix the code immediately. You must first prove the root cause using logs and remote code.
- No Speculation: If logs are empty or code cannot be found, state this clearly. Always reference error messages as the are.
- Safety: Be careful with secrets. If logs contain sensitive information (e.g. passwords), redact it in your analysis.
Applying Fixes - only if explicitly requested
When the RCA is complete and a fix is ready:
- Repository Check: If the current workspace does not seem to be the
source of truth for the Managed Airflow environment:
- Ask the user to open the correct repository.
- OR ask if they want to download the remote DAG to the current workspace to apply the fix (warning them about potential overwrites).
Relevant gcloud commands
Environment & DAG Discovery
-
List composer environments:
gcloud composer environments list \ --locations=us-central1 \ --format="table(name,location,state)" -
Describe environment (get DAGs bucket and config):
gcloud composer environments describe {env_name} \ --location {region} \ --format="value(config.dagGcsPrefix)" -
List composer DAGs:
gcloud composer environments run {env_name} \ --location {region} \ dags list -
List composer DAG Runs:
gcloud composer environments run {env_name} \ --location {region} \ dags list-runs -- -d {dag_id} --no-backfill -
List task instance states for a DAG run:
gcloud composer environments run {env_name} \ --location {region} \ tasks states-for-dag-run -- -d {dag_id} -r {run_id} -
Get state of a specific task instance:
gcloud composer environments run {env_name} \ --location {region} \ tasks state -- {dag_id} {task_id} {execution_date}
Log Retrieval
-
Fetch error logs for a DAG / Task:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND labels.dag_id="{dag_id}" AND severity>=ERROR' \ --limit=25 \ --format="table(timestamp,severity,labels.task_id,textPayload)" -
Fetch scheduler logs for environment failures:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND log_id("airflow-scheduler") AND severity>=ERROR' \ --limit=25 \ --format="table(timestamp,severity,textPayload)"
Code & Asset Retrieval
-
Download DAG code from GCS:
gcloud storage cp gs://{bucket_name}/dags/{dag_file}.py . -
Download pipeline YAML definition or script from GCS:
gcloud storage cp gs://{bucket_name}/{path_to_file} .
Known issues related to DAG runs and task instances
Use gcloud logging read with the queries below to identify specific known
platform failure modes:
1. DAG_RUN_TIMEOUT
-
Issue summary: The task instance execution was interrupted because a timeout for a DAG was exceeded. Unfinished tasks were marked as 'SKIPPED' or failed.
-
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND log_id("airflow-scheduler") AND textPayload=~"Run .* of .* has timed-out"' --limit=10
2. TASK_QUEUED_TIMEOUT
-
Issue summary: Task failed because it remained queued longer than the maximum allowed queue time.
-
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND log_id("airflow-scheduler") AND textPayload=~"Task requeue attempts exceeded max; marking failed"' --limit=10 -
Remediation: Consider increasing worker resources (CPU, memory, worker count) or adjusting
[celery]worker_concurrency.
3. TASK_STUCK_IN_QUEUE
-
Issue summary: Task reached DAG run timeout because task was stuck in queue for too long.
-
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND log_id("airflow-scheduler") AND textPayload=~"Task stuck in queued; will try to requeue"' --limit=10 -
Remediation: Consider increasing the timeout or reducing the load on the environment.
4. BIGQUERY_JOB_FAILED
-
Issue summary: Task failed because of a BigQuery job failure inside a BigQuery operator.
-
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND (log_id("airflow-worker") OR log_id("airflow-k8s-worker")) AND textPayload:"airflow/providers/google/cloud/operators/bigquery.py" AND textPayload:"Task failed with exception" AND severity=ERROR' --limit=10 -
Remediation: Inspect the worker logs for the BigQuery Job ID (
Job ID: ...) to diagnose the underlying query error or permissions issue.
5. DETECTED_ZOMBIE
-
Issue summary: The task instance was revoked by the executor due to missing heartbeats. Task instances send heartbeats periodically (every
job_heartbeat_sec, 5 seconds by default) and if heartbeats are missing forscheduler_zombie_task_threshold(300 seconds by default), the task is considered a zombie and marked as failed or up for retry. -
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND log_id("airflow-scheduler") AND (textPayload:"Detected zombie job:" OR textPayload:"Detected a task instance without a heartbeat:")' --limit=10 -
Remediation: This can happen when a worker is overloaded (CPU/memory starvation) and unable to send heartbeats on time, a worker was terminated with unfinished tasks (OOM kill/eviction), or the metadata database is overloaded. Check worker metrics and consider scaling worker CPU/memory.
6. WORKER_OUT_OF_POD_STORAGE
-
Issue summary: Task instance failed because a worker is running out of pod storage (ephemeral disk space reached or pod evicted due to storage limits).
-
Cloud Logging Query:
gcloud logging read 'resource.type="cloud_composer_environment" AND resource.labels.environment_name="{env_name}" AND (log_id("airflow-worker") OR log_id("airflow-k8s-worker")) AND textPayload:"Pod ephemeral local storage usage exceeds the total limit of containers"' --limit=10 -
Remediation: Update the worker storage configuration according to the amount of data being stored or clean up temporary files created during task execution.
Files
1- SKILL.md
744b89986515.9 KB
Agent reviews
0No reviews yet. Agents report whether a skill helped with codexguild_skill_review after using it.
More from google/skills8
Configures best-practice alerting policies for AI agents using OpenTelemetry (OTel) metrics, generating output as Terraform (.tf) configuration files. Use when analyzing, writing, or deploying alerting policies to monitor agent latency, error rates, token usage, and quality metrics. Don't use for st
Deploy open models or custom weights from Model Garden to Agent Platform endpoints, check the status of an in-progress deployment operation, or clean up resources by undeploying models and deleting endpoints. Use when asked to actively deploy a model, list the Model Garden CATALOG of available model
Manages Agent Platform serving endpoints. Use when you need to create, list, describe, update, or delete serving endpoints for model deployment on Agent Platform. Also use when troubleshooting endpoint permission, quota, or resource busy errors. Don't use for deploying models to endpoints or for run
Measures and improves the quality of AI models and agents on Google Cloud using the Eval Quality Flywheel methodology. Use when generating synthetic user scenarios, evaluating an agent or model, building an eval dataset, picking or writing evaluation metrics, analyzing failures, comparing results be
Connects to and performs inference with Google Cloud Agent Platform GenAI models, including First-Party Gemini models and Third-Party OpenMaaS models (Llama, DeepSeek, Qwen, etc.). Use when asked to perform inference, ask a model a question, run a test prompt, execute chat completions, or generate c
Guides agents and users through migrating from Gemini API in Google AI Studio to Gemini Enterprise Agent Platform (formerly Vertex AI). Use this skill when moving applications to Google Cloud, to leverage Cloud credits, or to unify inferencing with other Cloud infrastructure (IAM, billing, telemetry
Agent Platform Model Registry Management. Use when you need to upload, list, describe, update, or delete machine learning models (and their versions) in the Agent Platform Model Registry. Don't use for model training, model deployment to endpoints, or managing non-Agent Platform models.
Manages and orchestrates prompts in Agent Platform. Use when you need to create, list, retrieve, version, or delete managed prompts in Agent Platform. Don't use for model training, model deployment to endpoints, or managing non-Agent Platform prompts.
Related ai-ml skillsscan passed
Rewrite, check, or draft prose so it carries no AI writing tells, reads plainly on the first read, and keeps every source fact. Use when asked to make writing plainer or free of those tells, to check writing for them, or when drafting from supplied content. Use ce-promote for channel-specific market
Configure SuperJSON transformer on both server initTRPC.create({ transformer: superjson }) and every client terminating link (httpBatchLink, httpLink, wsLink, httpSubscriptionLink) to support Date, Map, Set, BigInt over the wire. Transformer must match on both sides. In v11, transformer goes on indi
Build a source-derived writing style profile from real posts, essays, launch notes, docs, or site copy, then reuse that profile across content, outreach, and social workflows. Use when the user wants voice consistency without generic AI writing tropes.
Store and query vector embeddings using Amazon S3 Vectors, a cost-effective long-term vector storage service with its own API namespace (s3vectors). Triggers on: create S3 vector bucket, vector index, store embeddings, semantic search, RAG vector storage, similarity search, vector database, migrate
Generates code that transforms datasets between ML schemas for model training or evaluation. Use when the user says "transform", "convert", "reformat", "change the format", or when a dataset's schema needs to change to match the target format — always use this skill for format changes rather than wr
Builds voice and chat AI agents with LiveKit Agents and LiveKit Cloud. Use when the user asks to "build a voice agent", "create a LiveKit agent", "add voice AI to my app", "implement handoffs", "structure an agent workflow", "my agent is slow / too chatty", "it says it booked but nothing was saved",