옵션 2: AWS Lambda로 Snowpipe 자동화

옵션 2: AWS Lambda로 Snowpipe 자동화

AWS Lambda는 이벤트로 트리거될 때 실행되어 시스템에 로드된 코드를 실행하는 컴퓨팅 서비스예요. 이 주제에서 제공하는 샘플 Python 코드를 적용하고 Snowpipe REST API를 호출해 외부 스테이지(즉, S3 버킷; Azure 컨테이너는 지원되지 않음)에서 데이터를 로드하는 Lambda 함수를 만들 수 있어요. 이 함수는 호스팅되는 AWS 계정에 배포돼요. Lambda에서 정의한 이벤트(예: S3 버킷의 파일이 업데이트될 때)가 Lambda 함수를 호출하고 Python 코드를 실행해요.

이 주제는 Snowpipe를 사용해 마이크로 배치 데이터를 연속적으로 자동 로드하도록 Lambda 함수를 구성하는 데 필요한 단계를 설명해요.

참고 — 이 주제는 Snowpipe REST API를 사용한 데이터 로드 준비의 지침으로 Snowpipe를 구성했다고 가정해요.

출처: Documentation

본문

1단계: Snowpipe REST API를 호출하는 Python 코드 작성

샘플 Python 코드

from __future__ import print_function
from snowflake.ingest import SimpleIngestManager
from snowflake.ingest import StagedFile
from requests import HTTPError
from cryptography.hazmat.primitives import serialization
from cryptography.hazmat.primitives.serialization import load_pem_private_key
from cryptography.hazmat.primitives.serialization import Encoding
from cryptography.hazmat.primitives.serialization import PrivateFormat
from cryptography.hazmat.primitives.serialization import NoEncryption
from cryptography.hazmat.backends import default_backend

import os

with open("./rsa_key.p8", 'rb') as pem_in:
  pemlines = pem_in.read()
  private_key_obj = load_pem_private_key(pemlines,
  os.environ['PRIVATE_KEY_PASSPHRASE'].encode(),
  default_backend())

private_key_text = private_key_obj.private_bytes(
  Encoding.PEM, PrivateFormat.PKCS8, NoEncryption()).decode('utf-8')
# Assume the public key has been registered in Snowflake:
# private key in PEM format

# List of files in the stage specified in the pipe definition
ingest_manager = SimpleIngestManager(account='<account_identifier>',
                   host='<account_identifier>.snowflakecomputing.com',
                   user='<user_login_name>',
                   pipe='<db_name>.<schema_name>.<pipe_name>',
                   private_key=private_key_text)

def handler(event, context):
  for record in event['Records']:
    bucket = record['s3']['bucket']['name']
    key = record['s3']['object']['key']

    print("Bucket: " + bucket + " Key: " + key)
    # List of files in the stage specified in the pipe definition
    # wrapped into a class
    staged_file_list = []
    staged_file_list.append(StagedFile(key, None))

    print('Pushing file list to ingest REST API')
    resp = ingest_manager.ingest_files(staged_file_list)

참고 — 샘플 코드는 오류 처리를 고려하지 않아요. 예를 들어 실패한 ingest_manager 호출을 재시도하지 않아요.

샘플 코드를 사용하기 전에 다음을 변경해요.

  • 보안 매개 변수 업데이트: private_key=""" / [REDACTED PRIVATE KEY] """ — 키 페어 인증 및 키 교체 사용에서 만든 개인 키 파일의 내용을 지정해요. PRIVATE_KEY_PASSPHRASE 환경 변수를 사용해 개인 키 파일을 복호화할 암호를 지정해요.
    • Linux 또는 macOS: export PRIVATE_KEY_PASSPHRASE='<passphrase>'
    • Windows: set PRIVATE_KEY_PASSPHRASE='<passphrase>'
  • 세션 매개 변수 업데이트:
    • account='<account_identifier>' — 계정의 고유 식별자(Snowflake가 제공)를 지정해요. host 설명을 참고해요.
    • host='<account_identifier>.snowflakecomputing.com' — Snowflake 계정의 고유 호스트 이름을 지정해요. 계정 식별자의 권장 형식은 *organization_name*-*account_name*이며, Snowflake 조직과 계정의 이름이에요. 자세한 내용은 조직의 계정 이름(형식 1, 권장)을 참고해요. 또는 필요하면 리전과 클라우드 플랫폼과 함께 계정 로케이터를 지정해요. 자세한 내용은 리전의 계정 로케이터(형식 2)를 참고해요.
    • user='<user_login_name>' — Snowpipe 코드를 실행할 Snowflake 사용자의 로그인 이름을 지정해요.
    • pipe='<db_name>.<schema_name>.<pipe_name>' — 데이터 로드에 사용할 파이프의 정규화된 이름을 <db_name>.<schema_name>.<pipe_name> 형식으로 지정해요.
  • 파일 객체 목록에서 가져올 파일 경로 지정: staged_file_list = [] — 지정하는 경로는 파일이 있는 스테이지에 상대적이어야 해요. 각 파일의 확장자를 포함한 전체 이름을 포함해요. 예를 들어 gzip으로 압축된 CSV 파일은 .csv.gz 확장자를 가질 수 있어요.
  • 파일을 편리한 위치에 저장해요. 이 주제의 나머지 지침은 파일 이름이 SnowpipeLambdaCode.py라고 가정해요.

2단계: Lambda 함수 배포 패키지 만들기

다음 지침을 완료해 Lambda용 Python 런타임 환경을 구축하고 (이 주제의) 1단계: Snowpipe REST API를 호출하는 Python 코드 작성에서 적용한 Snowpipe 코드를 추가해요. 이러한 단계에 대한 자세한 내용은 AWS Lambda 배포 패키지 문서를 참고해요(Python 지침 참고).

중요 — 다음 단계의 스크립트는 대표적인 예제이며, RPM에 의존하는 YUM 패키지 관리자를 사용하는 Amazon Machine Instance(AMI)를 기반으로 AWS EC2 Linux 인스턴스를 만든다고 가정해요. Debian 기반 Linux AMI를 선택하면 스크립트를 그에 맞게 업데이트해요.

  1. AWS EC2 지침을 완료해 AWS EC2 Linux 인스턴스를 만들어요. 이 인스턴스는 Snowpipe 코드를 실행할 컴퓨팅 리소스를 제공해요.
  2. SCP(Secure Copy)를 사용해 Snowpipe 코드 파일을 새 AWS EC2 인스턴스에 복사해요.
scp -i key.pem /<path>/SnowpipeLambdaCode.py ec2-user@<machine>.<region_id>.compute.amazonaws.com:~/SnowpipeLambdaCode.py

여기서:

  • <path>는 로컬 SnowpipeLambdaCode.py 파일의 경로.
  • <machine>.<region_id>는 EC2 인스턴스의 DNS 이름(예: ec2-54-244-54-199.us-west-2.compute.amazonaws.com). DNS 이름은 Amazon EC2 콘솔의 Instances 화면에 표시돼요.
  1. SSH(Secure SHell)를 사용해 EC2 인스턴스에 연결해요.
ssh -i key.pem ec2-user@<machine>.<region_id>.compute.amazonaws.com
  1. EC2 인스턴스에 Python과 관련 라이브러리를 설치해요.
sudo yum install -y gcc zlib zlib-devel openssl openssl-devel
wget https://www.python.org/ftp/python/3.6.1/Python-3.6.1.tgz
tar -xzvf Python-3.6.1.tgz
cd Python-3.6.1 && ./configure && make
sudo make install
sudo /usr/local/bin/pip3 install virtualenv
/usr/local/bin/virtualenv ~/shrink_venv
source ~/shrink_venv/bin/activate
pip install Pillow
pip install boto3
pip install requests
pip install snowflake-ingest
  1. .zip 배포 패키지(Snowpipe.zip)를 만들어요.
cd $VIRTUAL_ENV/lib/python3.6/site-packages
zip -r9 ~/Snowpipe.zip .
cd ~
zip -g Snowpipe.zip SnowpipeLambdaCode.py

3단계: Lambda용 AWS IAM 역할 만들기

AWS Lambda 문서를 따라 Lambda 함수를 실행할 IAM 역할을 만들어요.

역할의 IAM Amazon 리소스 이름(ARN)을 기록해요. 다음 단계에서 사용할 거예요.

4단계: Lambda 함수 만들기

(이 주제의) 2단계: Lambda 함수 배포 패키지 만들기에서 만든 .zip 배포 패키지를 업로드해 Lambda 함수를 만들어요.

aws lambda create-function \
--region us-west-2 \
--function-name IngestFile \
--zip-file fileb://~/Snowpipe.zip \
--role arn:aws:iam::<aws_account_id>:role/lambda-s3-execution-role \
--handler SnowpipeLambdaCode.handler \
--runtime python3.6 \
--profile adminuser \
--timeout 10 \
--memory-size 1024

--role에 대해 (이 주제의) 3단계: Lambda용 AWS IAM 역할 만들기에서 기록한 역할 ARN을 지정해요.

출력에서 새 함수의 ARN을 기록해요. 다음 단계에서 사용할 거예요.

5단계: Lambda 함수 호출 허용

S3에 함수를 호출하는 데 필요한 권한을 부여해요.

--source-arn에 대해 (이 주제의) 4단계: Lambda 함수 만들기에서 기록한 함수 ARN을 지정해요.

aws lambda add-permission \
--function-name IngestFile \
--region us-west-2 \
--statement-id enable-ingest-calls \
--action "lambda:InvokeFunction" \
--principal s3.amazonaws.com \
--source-arn arn:aws:s3:::<SourceBucket> \
--source-account <aws_account_id> \
--profile adminuser

6단계: Lambda 알림 이벤트 등록

Amazon S3 이벤트 알림 지침을 완료해 Lambda 알림 이벤트를 등록해요. 입력 필드에 (이 주제의) 4단계: Lambda 함수 만들기에서 기록한 함수 ARN을 지정해요.

더 알아보기