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

如何使用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:

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 .pem file (including the BEGIN/END header/footer lines) as the secret value. Name it something descriptive like sftp-production-private-key.
  • Ensure your Cloud Composer environment’s service account has the roles/secretmanager.secretAccessor permission 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 10:07:12