You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Airflow获取Dataproc作业ID并将Dataproc日志下载至Google Cloud Storage?

Great to hear you've got Airflow running Spark jobs on Dataproc already! Let's tackle your two requirements: capturing the Dataproc job ID and exporting those logs to a GCS bucket. Here's how to do both seamlessly within your Airflow workflow:

Solution Overview

We'll break this into two actionable steps:

  • Extracting the Dataproc job ID using Airflow's XCom mechanism
  • Configuring automatic log routing to GCS (or exporting existing logs if needed)

1. Capture Dataproc Job ID in Airflow

When you use Airflow's DataprocSubmitJobOperator to launch a Dataproc job, the operator automatically returns the unique job ID after successful submission. You can capture this ID using Airflow's XCom feature to reuse it in downstream tasks (like log exports or status checks).

Example Implementation

from airflow import DAG
from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator
from airflow.operators.python import PythonOperator
from datetime import datetime

def get_and_store_job_id(**context):
    # Pull the job ID from XCom pushed by the DataprocSubmitJobOperator
    job_id = context['task_instance'].xcom_pull(task_ids='submit_spark_job')
    print(f"Captured Dataproc Job ID: {job_id}")
    # Push the ID to a dedicated XCom key for easy downstream access
    context['task_instance'].xcom_push(key='dataproc_job_id', value=job_id)

with DAG(
    'dataproc_spark_workflow',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:

    submit_spark_job = DataprocSubmitJobOperator(
        task_id='submit_spark_job',
        project_id='your-gcp-project-id',
        region='your-gcp-region',
        job={
            "reference": {"project_id": "your-gcp-project-id"},
            "placement": {"cluster_name": "your-dataproc-cluster-name"},
            "spark_job": {
                "main_class": "com.yourorg.YourSparkJob",
                "jar_file_uris": ["gs://your-bucket/path/to/your-spark-job.jar"]
            }
        }
    )

    capture_job_id = PythonOperator(
        task_id='capture_job_id',
        python_callable=get_and_store_job_id,
        provide_context=True
    )

    submit_spark_job >> capture_job_id

The DataprocSubmitJobOperator pushes the job ID to XCom by default, so the downstream Python task can easily pull and store it for later use.

2. Export Dataproc Job Logs to GCS

You have two reliable options to get your logs into GCS, depending on whether you want to set this up upfront or after the job runs:

Option 1: Auto-Route Logs to GCS During Job Submission

This is the most efficient approach—configure your Dataproc job to stream logs directly to GCS as it runs. Just add a logging_config section to your job definition:

submit_spark_job = DataprocSubmitJobOperator(
    task_id='submit_spark_job',
    project_id='your-gcp-project-id',
    region='your-gcp-region',
    job={
        "reference": {"project_id": "your-gcp-project-id"},
        "placement": {"cluster_name": "your-dataproc-cluster-name"},
        "spark_job": {
            "main_class": "com.yourorg.YourSparkJob",
            "jar_file_uris": ["gs://your-bucket/path/to/your-spark-job.jar"]
        },
        # Add this to automatically send logs to GCS
        "logging_config": {
            "gcs_logging_config": {
                "bucket": "your-log-storage-bucket",
                "path": "dataproc-job-logs/{{ ds }}/{{ task_instance.task_id }}/"
            }
        }
    }
)

Logs will be organized in GCS using Airflow macros for date and task ID, making them easy to find later.

Option 2: Export Logs for Completed Jobs

If you need to export logs for a job that's already finished, use the gcloud CLI in an Airflow BashOperator, leveraging the job ID captured earlier:

from airflow.operators.bash import BashOperator

export_completed_logs = BashOperator(
    task_id='export_logs_to_gcs',
    bash_command="""
        gcloud dataproc jobs wait {{ ti.xcom_pull(key='dataproc_job_id') }} \
            --project=your-gcp-project-id \
            --region=your-gcp-region \
            --export-logs=gs://your-log-storage-bucket/retrieved-logs/{{ ti.xcom_pull(key='dataproc_job_id') }}/
    """
)

capture_job_id >> export_completed_logs

Critical Permissions Note

Ensure the Airflow service account has these roles:

  • roles/dataproc.editor (to submit jobs and access job metadata)
  • roles/storage.objectCreator and roles/storage.objectViewer on your target GCS bucket

内容的提问来源于stack exchange,提问作者Khilesh Chauhan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 15:17:42