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

新手求助:搭建Snowflake到Snowflake端到端数据架构及工具联动方案

端到端Snowflake-Snowflake数据架构实现方案

一、统一工作流核心:用Apache Airflow编排全流程

放弃拆分式触发逻辑,直接用Airflow作为唯一编排工具,串联Snowflake导出→SageMaker处理→S3转存→Snowflake加载全环节,确保工作流统一、状态可追踪。

二、各环节具体实现步骤

1. Snowflake数据导出到S3

直接用Snowflake的COPY INTO命令完成关联查询结果的导出,无需额外工具:

COPY INTO 's3://your-source-bucket/data/raw/'
FROM (
    SELECT t1.user_id, t2.order_amount, t3.product_category
    FROM user_table t1
    JOIN order_table t2 ON t1.user_id = t2.user_id
    JOIN product_table t3 ON t2.product_id = t3.product_id
)
STORAGE_INTEGRATION = s3_snowflake_integration
FILE_FORMAT = (TYPE = PARQUET COMPRESSION = SNAPPY);
  • 提前在Snowflake中创建STORAGE_INTEGRATION对象,赋予Snowflake读写目标S3桶的权限,避免硬编码密钥。

2. 基于已有Python代码创建SageMaker Pipeline并集成到Airflow

第一步:适配已有代码为SageMaker兼容脚本

确保代码支持从指定S3路径读取、处理后写入目标S3桶,示例框架:

import boto3
import pandas as pd

def process_data():
    s3_client = boto3.client('s3')
    # 读取S3原始数据
    obj = s3_client.get_object(Bucket='your-source-bucket', Key='data/raw/0001.parquet')
    raw_df = pd.read_parquet(obj['Body'])
    
    # 替换为你的已有处理逻辑
    processed_df = raw_df[raw_df['order_amount'] > 100].assign(amount_level='high')
    
    # 写入目标S3桶
    processed_df.to_parquet('s3://your-target-bucket/data/processed/output.parquet')

if __name__ == '__main__':
    process_data()
  • 将代码和依赖包(如pandas、pyarrow)整理成requirements.txt,一起上传到S3的代码存储桶(如s3://your-code-bucket/scripts/)。

第二步:定义SageMaker Pipeline

用SageMaker Python SDK创建包含处理步骤的Pipeline:

import sagemaker
from sagemaker.processing import ProcessingInput, ProcessingOutput, ScriptProcessor
from sagemaker.workflow.pipeline import Pipeline
from sagemaker.workflow.steps import ProcessingStep

sagemaker_session = sagemaker.Session()
execution_role = 'arn:aws:iam::123456789012:role/SageMakerExecutionRole'

# 定义脚本处理器,指定运行环境
script_processor = ScriptProcessor(
    command=['python3'],
    image_uri=sagemaker.image_uris.retrieve('sklearn', sagemaker_session.boto_region_name, 'latest'),
    role=execution_role,
    instance_count=1,
    instance_type='ml.t3.medium'
)

# 定义处理步骤
processing_step = ProcessingStep(
    name='DataCleaningAndEnrichment',
    processor=script_processor,
    inputs=[
        ProcessingInput(source='s3://your-source-bucket/data/raw/', destination='/opt/ml/processing/input'),
        ProcessingInput(source='s3://your-code-bucket/scripts/processing_script.py', destination='/opt/ml/processing/code'),
        ProcessingInput(source='s3://your-code-bucket/scripts/requirements.txt', destination='/opt/ml/processing/code')
    ],
    outputs=[
        ProcessingOutput(source='/opt/ml/processing/output', destination='s3://your-target-bucket/data/processed/')
    ],
    code='/opt/ml/processing/code/processing_script.py',
    arguments=['--requirements', '/opt/ml/processing/code/requirements.txt']
)

# 组装并发布Pipeline
pipeline = Pipeline(
    name='SnowflakeDataProcessingPipeline',
    parameters=[],
    steps=[processing_step],
    sagemaker_session=sagemaker_session
)

pipeline.upsert(role_arn=execution_role)

第三步:Airflow中触发SageMaker Pipeline

用Airflow的SageMakerPipelineOperator直接调用已创建的Pipeline:

from airflow.providers.amazon.aws.operators.sagemaker import SageMakerPipelineOperator

trigger_sagemaker_pipeline = SageMakerPipelineOperator(
    task_id='trigger_sagemaker_processing',
    pipeline_name='SnowflakeDataProcessingPipeline',
    aws_conn_id='aws_default',
    region_name='us-east-1',
    wait_for_completion=True
)

3. S3处理后数据加载回Snowflake

用Airflow的SnowflakeOperator执行COPY INTO命令完成加载:

from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator

load_to_snowflake = SnowflakeOperator(
    task_id='load_processed_data_to_snowflake',
    sql="""
        COPY INTO your_target_snowflake_table
        FROM 's3://your-target-bucket/data/processed/'
        STORAGE_INTEGRATION = s3_snowflake_integration
        FILE_FORMAT = (TYPE = PARQUET COMPRESSION = SNAPPY)
        ON_ERROR = 'CONTINUE';
    """,
    snowflake_conn_id='snowflake_default'
)

三、完整Airflow DAG编排

将所有任务按依赖关系串联,形成完整工作流:

from airflow import DAG
from airflow.utils.dates import days_ago
from airflow.providers.amazon.aws.operators.sagemaker import SageMakerPipelineOperator
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator

default_args = {
    'owner': 'data_team',
    'start_date': days_ago(1),
    'retries': 1
}

with DAG(
    'snowflake_end_to_end_pipeline',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False
) as dag:

    export_task = SnowflakeOperator(
        task_id='export_from_snowflake',
        sql="""
            COPY INTO 's3://your-source-bucket/data/raw/'
            FROM (
                SELECT t1.user_id, t2.order_amount, t3.product_category
                FROM user_table t1
                JOIN order_table t2 ON t1.user_id = t2.user_id
                JOIN product_table t3 ON t2.product_id = t3.product_id
            )
            STORAGE_INTEGRATION = s3_snowflake_integration
            FILE_FORMAT = (TYPE = PARQUET COMPRESSION = SNAPPY);
        """,
        snowflake_conn_id='snowflake_default'
    )

    process_task = SageMakerPipelineOperator(
        task_id='run_sagemaker_processing',
        pipeline_name='SnowflakeDataProcessingPipeline',
        aws_conn_id='aws_default',
        region_name='us-east-1',
        wait_for_completion=True
    )

    load_task = SnowflakeOperator(
        task_id='load_to_snowflake',
        sql="""
            COPY INTO your_target_snowflake_table
            FROM 's3://your-target-bucket/data/processed/'
            STORAGE_INTEGRATION = s3_snowflake_integration
            FILE_FORMAT = (TYPE = PARQUET COMPRESSION = SNAPPY)
            ON_ERROR = 'CONTINUE';
        """,
        snowflake_conn_id='snowflake_default'
    )

    # 定义任务执行顺序
    export_task >> process_task >> load_task

四、关键配置注意事项

  • 权限配置:
    • Snowflake存储集成需绑定S3桶的读写权限;
    • SageMaker执行角色需具备S3读写、SagePipeline执行权限;
    • Airflow的AWS/Snowflake连接需配置对应权限,确保能调用相关API和执行SQL。
  • 文件格式:优先使用Parquet列式存储,相比CSV占用空间更小、读写效率更高,Snowflake和SageMaker均有原生支持。
  • 监控与排障:通过Airflow UI查看全流程状态,通过SageMaker Studio查看Pipeline的处理日志,快速定位问题节点。

内容的提问来源于stack exchange,提问作者Logan27

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 07:48:15