如何通过Airflow将Snowflake表数据导出为CSV并上传至GCS
实现Snowflake数据导出CSV到GCS(Cloud Composer/Airflow)
完全可以通过Airflow直接完成需求,不需要替换成独立Python调度或Dataflow(当然这两者是备选方案,但Airflow原生就能高效实现)。下面给出两种最实用的实现方式:
方式一:Airflow PythonOperator中转导出(适合小/中数据量)
这种方式是在Airflow中查询Snowflake数据,转换为CSV后上传GCS,适合需要在导出过程中做简单数据处理的场景。
完整代码示例
from airflow import DAG from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook from airflow.providers.google.cloud.hooks.gcs import GoogleCloudStorageHook from airflow.operators.python import PythonOperator from datetime import datetime, timedelta import csv from io import StringIO def export_snowflake_to_gcs(**context): # 1. 连接Snowflake查询全量数据 dwh_hook = SnowflakeHook(snowflake_conn_id="snowflake_conn") # 注意:get_first仅返回第一行,改用get_records获取全部结果 records = dwh_hook.get_records("select col1,col2,col3,col4,col5 from table_name where col_name = '{{ prev_ds }}'") # 提取查询表头 cursor = dwh_hook.get_conn().cursor() cursor.execute("select col1,col2,col3,col4,col5 from table_name where col_name = '{{ prev_ds }}'") headers = [desc[0] for desc in cursor.description] # 2. 将数据转为CSV格式 output = StringIO() writer = csv.writer(output) writer.writerow(headers) writer.writerows(records) output.seek(0) # 3. 上传到GCS gcs_hook = GoogleCloudStorageHook(gcp_conn_id="google_cloud_default") gcs_bucket = "your-gcs-bucket-name" gcs_file_path = f"exports/snowflake_data_{{{{ ds }}}}.csv" gcs_hook.upload( bucket_name=gcs_bucket, object_name=gcs_file_path, data=output.getvalue(), mime_type="text/csv" ) # 定义DAG default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 1, 'retry_delay': timedelta(minutes=5) } with DAG( 'snowflake_to_gcs_csv', default_args=default_args, description='Export Snowflake data to CSV in GCS', schedule_interval='@daily', catchup=False ) as dag: export_task = PythonOperator( task_id='export_snowflake_to_gcs', python_callable=export_snowflake_to_gcs, provide_context=True ) export_task
关键注意点
- 替换
your-gcs-bucket-name为实际的GCS桶名 - 确保Airflow已配置好Snowflake连接(
snowflake_conn)和GCP连接(google_cloud_default) - 使用
get_records替代get_first获取全量数据,get_first仅返回查询结果的第一行,不适合导出全表数据 - 利用Airflow模板变量
{{ prev_ds }}和{{ ds }}处理日期逻辑,替代硬编码的日期函数
方式二:Snowflake COPY INTO直接导出(适合大数据量)
这种方式让Snowflake直接将数据导出到GCS,Airflow仅负责触发导出命令,性能更优,适合大数据量场景,不需要通过Airflow中转数据。
步骤1:配置Snowflake与GCS的集成
先在Snowflake中完成外部存储集成配置:
- 创建存储集成,指定GCS桶和授权服务账号
- 授予Snowflake角色访问该集成的权限
- 创建指向GCS桶的外部阶段
步骤2:Airflow中执行COPY INTO命令
from airflow import DAG from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 1, 'retry_delay': timedelta(minutes=5) } with DAG( 'snowflake_copy_to_gcs', default_args=default_args, description='Use Snowflake COPY INTO to export data to GCS', schedule_interval='@daily', catchup=False ) as dag: copy_task = SnowflakeOperator( task_id='snowflake_copy_to_gcs', snowflake_conn_id='snowflake_conn', sql=""" COPY INTO @gcs_external_stage/exports/snowflake_data_{{ ds }}.csv FROM (select col1,col2,col3,col4,col5 from table_name where col_name = '{{ prev_ds }}') FILE_FORMAT = (TYPE = CSV FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 0); """ ) copy_task
关键注意点
@gcs_external_stage是你在Snowflake中创建的指向GCS的外部阶段- 可通过Snowflake文件格式配置CSV的分隔符、引号、表头规则等参数
关于你考虑的两种方案的评估
- 方案1:用带调度的Python文件替代DAG:完全没必要。Airflow本身提供成熟的调度、监控、重试、日志管理能力,独立Python脚本需要自行实现这些功能,维护成本更高。
- 方案2:Dataflow(Beam):适合需要复杂数据转换、处理超大规模数据的场景。如果只是简单导出CSV,Airflow原生实现更轻量,不需要额外的Dataflow集群资源。
内容的提问来源于stack exchange,提问作者mkp24
相关产品推荐
相关产品推荐

