Snowpark 란

Snowpark는 Snowflake에서 제공하는 데이터 처리 프레임워크(라이브러리) 로, 데이터를 외부로 옮기지 않고 Snowflake 엔진 위에서 Python·Java·Scala 코드를 실행할 수 있게 합니다. Workspace Notebook은 Snowpark 코드를 작성·실행하는 개발 환경입니다.

차별점

1. 전통적인 데이터 처리 방식

과거에는 데이터를 분석하거나 가공하려면 데이터가 저장된 곳에서 데이터를 꺼내 별도의 처리 환경으로 옮겨야만 했습니다. 하지만 데이터의 양이 기하급수적으로 늘어난 오늘날, 이러한 '데이터 이동(Data Movement)' 방식은 데이터 분석의 걸림돌이 되고 있습니다.

전통적인방식의한계

  • 데이터가 무거워질수록 비용과 시간은 늘어납니다 (낮은 성능, 높은 비용)

    • 전통적인 방식에서는 분석을 위해 기가바이트(GB)에서 테라바이트(TB) 단위의 데이터를 네트워크를 통해 전송해야 합니다.
  • 전문가 없이는 돌아가지 않는 복잡한 시스템 (복잡한 튜닝)

    • 기존 환경은 데이터가 들어올 때마다 서버를 수동으로 늘리거나(Scaling), 데이터베이스 최적화 설정을 매번 바꿔줘야 했습니다.
    • 비즈니스 부서에서는 당장 데이터 분석 결과가 필요한데, IT 전문가가 서버 클러스터를 튜닝할 때까지 며칠을 기다려야 하는 상황이 빈번하게 발생합니다.
  • 데이터가 이동할수록 보안은 취약해집니다 (보안과 거버넌스)

    • 데이터가 저장소를 떠나 외부 환경으로 이동하는 순간, 그 데이터는 통제권을 벗어나게 됩니다.
    • "이 데이터가 어디로 복사되었지?", "누가 이 복사본을 보고 있지?"라는 질문에 답하기 어려워집니다. 이는 심각한 보안 사고나 데이터 거버넌스 위반으로 이어질 수 있습니다.

2. Snowpark의 새로운 처리 방식

전통적인 데이터 분석 환경에서는 Python이나 Java 코드를 실행하기 위해 데이터를 외부로 추출해야만 했습니다. 하지만 Snowpark는 이 패러다임을 완전히 뒤바꾸어 데이터가 있는 Snowflake 엔진 안에서 직접 코드를 실행하게 만듭니다.

Snowpark처리방식

  • 익숙한 언어로, 데이터가 있는 곳에서 직접 개발하세요

    • Snowpark는 개발자가 선호하는 프로그래밍 언어 (Python, Java, Scala)를 그대로 사용할 수 있게 해줍니다.
    • Snowpark API: 복잡한 SQL 문을 길게 작성할 필요 없이, Python의 DataFrame 형식으로 데이터 파이프라인과 앱을 구축할 수 있습니다.
    • 커스텀 코드 작성: UDF(사용자 정의 함수)나 프로시저를 통해 자신만의 비즈니스 로직을 Snowflake 내부에 직접 배포할 수 있습니다.
  • 머신러닝(ML)의 전 과정을 Snowflake 안에서 (Snowpark ML)

    • Data Scientist들에게 가장 반가운 소식은 ML 워크플로우를 위해 데이터를 외부 ML 서버로 옮길 필요가 없다는 점입니다.
    • End-to-End ML: 데이터 전처리(Feature Engineering)부터 모델 학습, 배포까지 모든 과정을 Snowflake의 강력한 인프라 위에서 안전하게 수행할 수 있습니다.
  • 유연하고 강력한 실행 환경 (CPU부터 GPU까지)

    • Snowpark는 단순한 라이브러리를 넘어, 코드가 돌아가는 최적의 '그릇'을 제공합니다.
    • Workspace Notebook에서는 Python 코드는 Container Runtime(Compute Pool) 에서 실행되고, Snowpark DataFrame 연산과 SQL은 가상 웨어하우스(Virtual Warehouse) 로 pushdown됩니다. GPU·ML 워크로드는 Compute Pool 기반 Container Runtime 에서 처리합니다
    • Snowpark Container Services (SPCS): 더 복잡한 로직이나 고성능 연산이 필요할 때, GPU를 포함한 사용자 정의 런타임 환경을 컨테이너 기반으로 제공하여 어떤 언어라도 실행할 수 있게 합니다.

예시 코드로 비교

1. 클라이언트 사이드 처리

  • 예시코드
    
    # Snowpark DataFrame을 pandas로 변환
    snowpark_df = session.table("my_table")
    pandas_df = snowpark_df.to_pandas()  # 모든 데이터를 로컬로 가져옴
    
    # 이후 모든 처리는 로컬 메모리에서 진행
    result = pandas_df.groupby('category').sum()
  • 특징
    • 즉시 실행(Eager Evaluation): 각 연산마다 즉시 결과를 메모리에 저장
    • 단일 머신 처리: 로컬 메모리에서만 작업
    • 메모리 제한: 머신의 물리적 메모리 한계에 제약
    • 데이터 이동: Snowflake → 로컬로 모든 데이터 전송

2. 서버 사이드 처리

  • 예시코드
    
    # Snowpark DataFrame으로 처리 (데이터는 Snowflake에 유지)
    snowpark_df = session.table("my_table")
    result = snowpark_df.group_by('category').sum('amount')  # SQL로 변환되어 실행
    result.show()  # 이 시점에서만 실제 실행
  • 특징
    • 지연 평가(Lazy Evaluation): 실제 액션이 호출될 때까지 실행 지연
    • 분산 컴퓨팅: Snowflake 클러스터의 모든 노드 활용
    • 무제한 확장: 메모리 제한 없이 대용량 데이터 처리
    • 데이터 유지: 데이터가 Snowflake에서 이동하지 않음

적용 가이드

측면클라이언트 사이드 처리
to_pandas() 사용
서버 사이드 처리
Snowpark (인-데이터베이스 처리)
데이터 크기메모리 제한 (GB급)무제한 (TB, PB급)
처리 속도소량 데이터에 빠름대량 데이터에 빠름
네트워크 비용전체 데이터 전송결과만 전송
스케일링수직 확장만 가능수평 확장 가능
최적화pandas 최적화SQL 엔진 최적화
사용 가이드데이터가 수십만 행 이하로 작을 때
기존 pandas 코드를 그대로 사용해야 할 때
로컬 시각화나 분석 라이브러리 사용 시
즉시 피드백이 필요한 탐색적 분석
데이터가 수백만 행 이상으로 클 때
ETL/ELT 파이프라인 구축 시
성능과 비용 효율성이 중요할 때
프로덕션 워크로드 개발 시
주의사항메모리 모니터링 필수 (OOM 에러 방지)
필터링 먼저: to_pandas() 전에 데이터 양 줄이기
배치 처리: 청크 단위로 나누어 처리
지연 평가 이해: 실제 실행 시점 파악
액션 메서드 활용: .show(), .collect(), .count()
데이터 타입 매핑 주의

Snowpark Dataframe 실습

실습 전제

  • 본 Snowpark 실습은 Notebooks in Workspaces 환경에서 진행합니다.
  • Legacy Snowflake Notebooks가 아니라, Workspace에 저장된 .ipynb 파일을 열어 코드를 실행하는 방식입니다.
  • 따라서 이전 Snowflake Notebooks 실습을 먼저 완료해 두는 것을 권장합니다.
    • setup.sql 실행으로 DS_ROLE, TESTDB, Compute Pool(DS_NOTEBOOK_POOL) 준비
    • Workspace에서 Notebook Service 생성 및 연결
    • SQL / Python / Markdown 셀 실행, Jinja 변수 참조 등 기본 Notebook 조작

1. 새 Notebook 생성

Snowpark 실습을 시작하기 전에, Workspace와 Notebook Service가 아래 조건을 만족하는지 확인하세요.

  1. Workspace에서 .ipynb 파일 생성
    Snowpark 실습용 노트북은 Legacy Notebooks 메뉴가 아니라 Projects → Workspaces에서 만듭니다.

    • Projects → Workspaces로 이동
    • 원하는 폴더에서 + 또는 + Add new 클릭
    • Notebook 선택 후 파일 이름을 Test_snowpark.ipynb로 입력
    • 생성된 .ipynb 파일을 열어 실습 진행

      Workspace 노트북은 database/schema에 자동으로 저장되지 않으므로, 이후 SQL·Snowpark 코드에서 테이블을 조회할 때는 fully qualified name을 사용하거나 USE DATABASE / USE SCHEMA로 실행 컨텍스트를 지정해야 합니다.

  2. Notebook Service 연결 (Connected 상태)
    Workspace Notebook에서 Python 코드를 실행하려면 Notebook Service가 필요합니다. Container Runtime 기반 Notebook Service는 Compute Pool 위에서 동작하며, 이전 Notebook 실습에서 생성한 노트북서비스를 재사용할 수 있습니다.

    • Test_Snowpark.ipynb 에서 상단 Connected 드롭다운을 클릭
    • 이전 실습에서 생성한 Notebook Service 선택
    • 상태가 Connected이면 노트북 셀 실행 가능

      새 Notebook Service가 필요하면 Create new service를 클릭합니다.

  3. Role / Warehouse 설정
    Workspace Notebook 상단의 Role / Warehouse picker에서 아래를 설정합니다.

    • Role: DS_ROLE
    • Warehouse: COMPUTE_WH (또는 실습용으로 준비한 Query Warehouse)

      SQL 쿼리와 Snowpark DataFrame의 pushdown 연산은 Query Warehouse에서 실행됩니다. Python 커널은 Notebook Service(Compute Pool) 에서 실행되므로, Role·Warehouse·Compute Pool 권한을 모두 갖춰야 실습이 정상적으로 진행됩니다.

  4. snowflake_sample_data 접근 권한
    본 실습은 Snowflake 샘플 데이터셋 snowflake_sample_data.tpch_sf10.orders를 사용합니다. DS_ROLE 또는 현재 사용 중인 Role에 해당 database/schema/table에 대한 USAGE 권한이 있어야 합니다.

    • 권한이 없으면 ACCOUNTADMIN 또는 관리자 Role로 아래와 같이 권한을 부여합니다.
      GRANT IMPORTED PRIVILEGES ON DATABASE snowflake_sample_data TO ROLE DS_ROLE;

2. SQL셀 데이터를 Python으로

첫번째 셀을 SQL셀로 변경하여 다음 쿼리를 실행합니다

  • 셀 이름 : S1
  • 다음 내용을 복사해서 넣고 실행합니다
    
    select * 
      from snowflake_sample_data.tpch_sf10.orders
     limit 1000;

S1

Python셀을 추가하여 다음 코드를 실행합니다

  • 셀 이름 : P1
  • 다음 내용을 복사해서 넣고 실행합니다
    
    df = dataframe_1
    dt = df.query('O_TOTALPRICE > 200000')
    dt.head()

P1

  • SQL 셀 S1 실행 후, Python 셀에서 결과를 참조할 때는 셀 이름 S1 대신 결과 tooltip에 표시된 dataframe_x 이름을 사용하세요.

3. DataFrame 속도 비교

  • Python 셀을 추가하여 데이터베이스 연결 세션을 정의합니다

    
    from snowflake.snowpark.context import get_active_session
    session = get_active_session()
  • Python 셀을 추가하여 to_pandas() DataFrame으로 데이터를 조회합니다

    
    order_df = session.table("snowflake_sample_data.tpch_sf10.orders").to_pandas();
    
    customer_revenue = order_df.groupby('O_CUSTKEY')['O_TOTALPRICE'].sum().reset_index()
    
    # 매출 내림차순으로 정렬하여 상위 10명 추출
    top_10_customers = customer_revenue.sort_values(by='O_TOTALPRICE',ascending=False).head(10)
    
    # 결과 출력
    top_10_customers.head(10)
    • tpch_sf10.orders는 매우 큰 테이블입니다. LIMIT 없이 to_pandas()하면 OOM·장시간 실행·높은 비용이 날 수 있습니다.
    • 본 실습은 의도적으로 대용량 비교이므로 실제 사용시에는 LIMIT 100000 등 상한 추가하는 것이 바람직합니다

to_pandas

  • 다시 Python 셀을 추가하여 Snowpark DataFrame으로 데이터를 조회합니다

    
    from snowflake.snowpark.functions import sum as snowflake_sum, col
    
    df_snowpark = session.table("snowflake_sample_data.tpch_sf10.orders")
    
    df_customer_revenue = df_snowpark.group_by(col("O_CUSTKEY")).agg(snowflake_sum(col("O_TOTALPRICE")).alias("O_TOTALPRICE"))
    
    # 매출 내림차순 정렬 및 상위 10명 추출
    df_top_10_customers = df_customer_revenue.sort(col("O_TOTALPRICE"),ascending=False).limit(10)
    
    # Pandas DataFrame으로 변환
    df_pandas = df_top_10_customers.to_pandas()
    
    # 결과 출력
    df_pandas.head(10)
    • 셀 상단의 수행 시간을 비교해 보세요
    • 첫 실행·캐시·Virtual Warehouse 크기·Compute Pool 상태에 따라 차이가 큽니다.
    • Virtual Warehouse, Compute Pool 크기 변화에 따른 성능 테스트도 추가적으로 해 볼 수 있습니다.

      Workspace Notebook에서 to_pandas()는 Container Runtime 메모리로 데이터를 가져오며, Snowpark 연산의 pushdown 실행은 Query Warehouse에서 이루어집니다. 따라서 Compute Pool 크레딧과 Warehouse 크레딧이 함께 발생할 수 있습니다.

snowpark

profile
Snowflake Korea SE

0개의 댓글