방법1. Airflow S3 Hook 사용하기
방법2. Airflow S3 boto 사용하기


0. AWS CLI에 엑세스 키 등록


aws configure

AWS Access Key ID : {Access Key ID}
AWS Secret Access Key : {Secret Access Key}


Xixia


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 객체의 메타데이터 정보를 반환합니다.