如何通过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:
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.objectCreatorandroles/storage.objectVieweron your target GCS bucket
内容的提问来源于stack exchange,提问作者Khilesh Chauhan

