옵션 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>'
- Linux 또는 macOS:
- 세션 매개 변수 업데이트:
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를 선택하면 스크립트를 그에 맞게 업데이트해요.
- AWS EC2 지침을 완료해 AWS EC2 Linux 인스턴스를 만들어요. 이 인스턴스는 Snowpipe 코드를 실행할 컴퓨팅 리소스를 제공해요.
- 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 화면에 표시돼요.
- SSH(Secure SHell)를 사용해 EC2 인스턴스에 연결해요.
ssh -i key.pem ec2-user@<machine>.<region_id>.compute.amazonaws.com
- 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
- .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을 지정해요.