如何通过v1beta2 REST API在Airflow中创建自定义镜像Dataproc集群
I’ve helped folks work through this exact scenario before—here’s how you can replicate your gcloud beta dataproc clusters create command using the Dataproc v1beta2 REST API within Airflow:
Key Background
The --image flag in your gcloud command maps directly to the config.gceClusterConfig.imageUri field in the v1beta2 API request body. You’ll need to send a POST request to the Dataproc clusters create endpoint, with this field populated with your custom image name or full URI.
Option 1: Use Airflow’s HttpOperator
This is a straightforward approach for direct REST calls. You’ll need to set up authentication via Airflow’s Google Cloud connection or a service account token:
from airflow.providers.http.operators.http import HttpOperator from airflow.utils.dates import days_ago import json create_dataproc_cluster = HttpOperator( task_id="create_dataproc_cluster", method="POST", http_conn_id="google_cloud_dataproc", # Configure this connection in your Airflow UI endpoint="v1beta2/projects/your-project-id/regions/your-region/clusters", headers={"Content-Type": "application/json"}, data=json.dumps({ "clusterName": "your-cluster-name", "config": { "gceClusterConfig": { "imageUri": "custom-image-name", # Matches your gcloud --image flag "zoneUri": "your-target-zone" # e.g., us-central1-a }, "masterConfig": { "machineTypeUri": "n1-standard-2" }, "workerConfig": { "machineTypeUri": "n1-standard-2", "numInstances": 2 } # Add other cluster configs (like initialization actions) as needed } }), dag=your_dag_instance )
Option 2: Use PythonOperator with Google API Client
For more flexibility (like dynamic parameter handling or retry logic), use the google-api-python-client library to call the API directly. This integrates seamlessly with Airflow’s Google Cloud authentication:
from airflow.operators.python import PythonOperator from googleapiclient.discovery import build from google.auth import default import json def create_cluster_with_custom_image(): # Fetch default credentials (uses Airflow's GCP connection or attached service account) credentials, project_id = default() dataproc_client = build('dataproc', 'v1beta2', credentials=credentials) cluster_body = { "clusterName": "your-cluster-name", "config": { "gceClusterConfig": { "imageUri": "custom-image-name", "zoneUri": "your-target-zone" }, "masterConfig": { "machineTypeUri": "n1-standard-2" }, "workerConfig": { "machineTypeUri": "n1-standard-2", "numInstances": 2 } } } # Send the cluster creation request request = dataproc_client.projects().regions().clusters().create( projectId=project_id, region="your-region", body=cluster_body ) response = request.execute() print(f"Cluster creation started: {json.dumps(response)}") create_cluster_task = PythonOperator( task_id="create_dataproc_cluster", python_callable=create_cluster_with_custom_image, dag=your_dag_instance )
Important Notes
- Authentication: Ensure your Airflow environment has the
Dataproc Editorrole (or equivalent permissions) to create clusters. If using a service account, attach it to your Airflow workers or configure the GCP connection in the Airflow UI. - Full Image URI: If your custom image lives in a different project, use the full format:
projects/your-image-project/global/images/custom-image-name. - Region Consistency: Match the region in the API endpoint to your cluster’s zone to avoid unexpected errors.
内容的提问来源于stack exchange,提问作者Georges Kohnen

