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


0. Airflow S3 Hook 사용하기


전제조건
s3 생성 및 Access key 생성


1. 라이브러리 설치


pip3 install apache-airflow[amazon]


2. Airflow Connection 등록


Airflow UI에서 Admin -> connection 탭에 들어가 + 버튼을 클릭하여 새 연결을 설정해 줍니다.

Xixia

Connection Id : 사용할 ID
Connection Type : Amazon Web Services
AWS Access Key ID , AWS Secret Access Key
Extra : {"region_name : "ap-northeast-2"}

Xixia

3. s3 업로드 예시 코드

from airflow.decorators import dag, task, task_group
from datetime import datetime
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
import os

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():


    @task
    def upload_to_s3(**context):
        print('start')
        filename = f'/home/ec2-user/test.txt'
        key = f'data/jjw_test.txt'
        bucket_name = 'jjw-raw-bucket'

        # S3에 파일 업로드
        hook = S3Hook('aws_access_key')  # connection ID
        if os.path.exists(filename):
            hook.load_file(filename=filename, key=key, bucket_name=bucket_name, replace=True)
            print(f"File {filename} uploaded to S3 bucket {bucket_name} with key {key}")
        else:
            print(f"File {filename} not found.")


    upload_to_s3()

dag = test()


task 확인

Xixia


s3 확인

Xixia



3. s3 파일 읽기 예시 코드

from airflow.decorators import dag, task, task_group
from datetime import datetime
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
import os

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():


    @task
    def read_file_from_s3(**context):
        bucket_name = 'jjw-raw-bucket'
        key = f'data/jjw_test.txt'
        
        hook = S3Hook('aws_access_key')
        
        # S3에서 파일 읽기
        file_obj = hook.get_key(key, bucket_name=bucket_name)
        if file_obj:
            file_content = file_obj.get()['Body'].read().decode('utf-8')
            print(f"Read content from {key}:\n{file_content}")
            return file_content  # 파일 내용을 다음 태스크에 사용할 수 있도록 반환
        else:
            print(f"File {key} not found in bucket {bucket_name}.")
            return None


    read_file_from_s3()

dag = test()



3. s3 파일 삭제 예시 코드

from airflow.decorators import dag, task, task_group
from datetime import datetime
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
import os

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():


    @task
    def delete_file_from_s3(**context):
        bucket_name = 'khuda-de-project'
        key = f'data/jjw_test.txt'
        
        hook = S3Hook('de_velog_aws')
        
        # S3에서 파일 삭제
        hook.delete_objects(bucket=bucket_name, keys=key)
        print(f"File {key} deleted from S3 bucket {bucket_name}")


    delete_file_from_s3()

dag = test()