방법1. Airflow S3
Hook사용하기
방법2. Airflow S3boto사용하기
0. AWS CLI에 엑세스 키 등록
aws configure
AWS Access Key ID : {Access Key ID}
AWS Secret Access Key : {Secret Access Key}

1. s3 업로드, 읽기, 삭제 예시 코드
from airflow.decorators import dag, task, task_group
from datetime import datetime
import boto3
default_args = {
'owner': 'jinwoo',
'start_date': datetime(2024, 11, 10), # 시작 날짜
# 'end_date': datetime(2024, 11, 11), # 시작 날짜
}
@dag(
# schedule_interval="0 4 * * *", # UTC
schedule_interval=None,
default_args=default_args,
tags=['task_test']
)
def test():
file_path = f'/home/ec2-user/test.txt'
bucket_name = 'jjw-raw-bucket'
key = f'data/jjw_test.txt'
# 파일 업로드
@task
def upload_to_s3(**context):
print('start....')
s3_client = boto3.client('s3')
s3_client.upload_file(file_path, bucket_name, key)
print(f"File uploaded to {bucket_name}/{key}")
# 파일 읽기
@task
def read_from_s3(**context):
print('start....')
s3_client = boto3.client('s3')
response = s3_client.get_object(Bucket=bucket_name, Key=key)
file_content = response['Body'].read().decode('utf-8')
print(f"File content: \n{file_content}")
# 파일 삭제
@task
def delete_from_s3(**context):
print('start....')
s3_client = boto3.client('s3')
s3_client.delete_object(Bucket=bucket_name, Key=key)
print(f"File deleted from {bucket_name}/{key}")
upload_to_s3() >> read_from_s3() >> delete_from_s3()
dag = test()
2. boto3 S3 Client 주요 함수들
| 함수명 | 설명 |
|---|---|
| list_buckets() | 현재 계정에서 접근 가능한 S3 버킷들의 목록을 반환합니다. |
| create_bucket(Bucket) | 새로운 S3 버킷을 생성합니다. |
| delete_bucket(Bucket) | 지정된 S3 버킷을 삭제합니다. |
| list_objects_v2(Bucket, Prefix) | 특정 버킷의 객체 목록을 반환합니다. (Prefix를 사용하여 폴더처럼 특정 경로만 필터링 가능) |
| upload_file(Filename, Bucket, Key) | 로컬 파일을 S3 버킷의 지정된 키(Key)로 업로드합니다. |
| download_file(Bucket, Key, Filename) | S3 버킷에서 특정 키(Key)의 파일을 로컬로 다운로드합니다. |
| delete_object(Bucket, Key) | 특정 버킷에서 지정된 객체(Key)를 삭제합니다. |
| put_object(Bucket, Key, Body) | 지정된 객체(Key)로 데이터를 업로드합니다. (파일 대신 문자열이나 바이너리 데이터 업로드 가능) |
| get_object(Bucket, Key) | S3 객체의 데이터를 반환합니다. (Body를 사용하여 파일 내용을 읽을 수 있음) |
| copy(CopySource, Bucket, Key) | S3 내에서 객체를 복사합니다. |
| head_object(Bucket, Key) | S3 객체의 메타데이터 정보를 반환합니다. |