如何使用Cloud Composer下载及访问文件?两类场景技术问询
Great questions—let’s break down each scenario with Cloud Composer-specific best practices that align with your local Airflow workflow but fit the managed environment constraints:
1. Storing and Accessing SFTP Private Key (.pem) in Cloud Composer
Unlike your local Airflow setup where you can drop keys in a /keys folder alongside your DAGs, Cloud Composer is a fully managed environment—you don’t have persistent, direct access to worker node filesystems for storing sensitive credentials like .pem keys. Here are the most secure options:
Option 1: Google Cloud Secrets Manager (Recommended)
This is the gold standard for sensitive data, offering fine-grained access control, audit logging, and automatic rotation capabilities:
- Head to the Google Cloud Console, create a new secret in Secrets Manager, and paste the full content of your
.pemfile (including theBEGIN/ENDheader/footer lines) as the secret value. Name it something descriptive likesftp-production-private-key. - Ensure your Cloud Composer environment’s service account has the
roles/secretmanager.secretAccessorpermission for this secret (you can set this via IAM in the Console). - In your DAG, fetch the secret and use it with the
SFTPHook:
from airflow.providers.google.cloud.hooks.secret_manager import GoogleCloudSecretManagerHook from airflow.providers.sftp.hooks.sftp import SFTPHook def get_sftp_connection(): secret_hook = GoogleCloudSecretManagerHook() # Fetch secret and decode from bytes to string private_key = secret_hook.get_secret( secret_id="sftp-production-private-key", project_id="your-gcp-project-id" ).decode("utf-8") sftp_hook = SFTPHook( ssh_conn_id="your-sftp-airflow-connection", private_key=private_key ) return sftp_hook
Option 2: Embed in Airflow Connection Extra Field
For simpler use cases (less secure than Secrets Manager but still acceptable for non-production or low-risk environments):
- Go to your Airflow UI > Admin > Connections, create an SFTP-type connection.
- In the Extra field, paste your private key as a JSON string (escape newlines with
\n):
{ "private_key": "-----BEGIN RSA PRIVATE KEY-----\nMIIEpAIBAAKCAQEAz...\n-----END RSA PRIVATE KEY-----" }
- You can then reference this connection directly in SFTP operators without extra code:
from airflow.providers.sftp.operators.sftp import SFTPOperator fetch_sftp_file = SFTPOperator( task_id="fetch_sftp_file", ssh_conn_id="your-sftp-airflow-connection", remote_filepath="/remote/sftp/path/data.csv", local_filepath="/tmp/temp-data.csv" # Composer manages this temp path )
2. Migrating Files from SFTP to Cloud Storage
You don’t need to manually manage intermediate file storage like you do in local Airflow. Cloud Composer offers built-in tools to handle this seamlessly, either via a dedicated operator or a custom hook-based workflow:
Option 1: Use SFTPToGCSOperator (Simplest Approach)
This purpose-built operator handles the entire transfer process, including temporary storage on Composer workers (you don’t need to worry about filesystem paths):
from airflow.providers.google.cloud.transfers.sftp_to_gcs import SFTPToGCSOperator sftp_to_gcs_transfer = SFTPToGCSOperator( task_id="sftp_to_gcs", ssh_conn_id="your-sftp-airflow-connection", source_path="/remote/sftp/path/*.parquet", # Supports wildcards for bulk transfers destination_bucket="your-gcs-bucket-name", destination_path="processed-data/sftp-imports/", gcp_conn_id="google_cloud_default" # Uses Composer's default GCP service account )
Option 2: Custom Workflow with Hooks (For Advanced Logic)
If you need to modify files mid-transfer (e.g., parsing, filtering), use tempfile to create a temporary local file (auto-cleaned after execution) and transfer via hooks:
import tempfile from airflow.providers.sftp.hooks.sftp import SFTPHook from airflow.providers.google.cloud.hooks.gcs import GoogleCloudStorageHook from airflow.operators.python import PythonOperator def custom_sftp_to_gcs_transfer(): sftp_hook = SFTPHook(ssh_conn_id="your-sftp-airflow-connection") gcs_hook = GoogleCloudStorageHook(gcp_conn_id="google_cloud_default") # Download SFTP file to a temp file with tempfile.NamedTemporaryFile(mode="wb") as temp_file: sftp_hook.retrieve_file( remote_path="/remote/sftp/path/raw-data.csv", local_path=temp_file.name ) # Optional: Add custom processing here (e.g., edit temp_file) # Upload to GCS gcs_hook.upload( bucket_name="your-gcs-bucket-name", object_name="processed-data/cleaned-data.csv", filename=temp_file.name ) custom_transfer_task = PythonOperator( task_id="custom_sftp_to_gcs", python_callable=custom_sftp_to_gcs_transfer )
内容的提问来源于stack exchange,提问作者Aaron Deutsch

