Python 가상 로그 생성부터 AWS Firehose를 거쳐 S3 적재까지

정성헌·2026년 4월 15일

실시간 스트리밍 데이터를 다룰 때, 발생하는 모든 로그를 S3에 실시간으로 '직접' 꽂아 넣는 것은 안티 패턴(Anti-pattern)이다. S3에 자잘한 파일이 무수히 쌓이게 되면(Small File Problem), 추후 Athena나 Spark로 데이터를 읽어 들일 때 I/O 병목이 발생하고 탐색 비용이 폭발하기 때문이다.

이 문제를 해결하기 위해 중간에 버퍼(Buffer) 역할을 해주는 Amazon Data Firehose(ADF)를 배치한다. 오늘은 Python으로 가상의 로그를 생성하고, Firehose를 통해 일정한 조건(크기/시간)에 맞춰 S3에 덩어리(Batch) 단위로 적재하는 파이프라인을 구축해 본다.

1. 아키텍처 및 파이프라인 흐름

전체적인 데이터의 흐름은 다음과 같다.
1. Producer: Python Faker 라이브러리를 이용해 금융/이커머스 등 가상 도메인 로그 생성
2. Delivery: Boto3 (AWS SDK)를 통해 Firehose 스트림으로 데이터 전송 (put_record)
3. Buffering & Sink: Firehose가 지정된 버퍼 조건(예: 1MiB or 60초)을 채운 뒤 S3에 파일로 저장


2. 가짜 데이터 생성 및 직렬화 (Log Generator)

데이터 파이프라인 테스트를 위해 가장 먼저 할 일은 현실성 있는 데이터를 만드는 것이다. Faker 라이브러리를 활용해 금융(Finance) 트랜잭션 로그를 생성하는 클래스를 구성한다.

# log_generator.py (일부 발췌)
from faker import Faker
from datetime import datetime
import random

fake = Faker('ko_KR')

class LogGenerator:
    def finance(self):
        return {
            "timestamp" : datetime.now().isoformat(),
            "industry"  : "FINANCE",
            "user_id"   : fake.uuid4()[:12],
            "transaction" : random.choice(['이체','출금','입금','결제']),
            "amount" : round(random.uniform(1000, 1000000), -2),
            "status" : random.choices(["SUCCESS", "FAIL"], weights=[0.95, 0.05])[0]
        }

이 객체 형태의 데이터를 AWS로 전송하려면 문자열 형태(JSON str)로 변환하는 직렬화(Serialization) 과정이 필수적이다.

# run.py
import json

def make_one_log():
    # ensure_ascii=False 옵션으로 한글 깨짐 방지
    return json.dumps(log_gen.finance(), ensure_ascii=False) 

3. Boto3를 이용한 AWS Firehose 연결 (인증)

이제 Python 스크립트가 AWS 인프라에 접근하기 위한 '통행증'을 발급받아야 한다. boto3 라이브러리를 사용해 Firehose 클라이언트를 생성한다.

여기서 가장 중요한 것은 코드가 실행되는 환경에 따른 인증 방식의 분리다. 하드코딩된 자격 증명(Access Key)은 보안 사고의 주범이므로 절대 피해야 한다.

# adf_direct_data_put.py
import boto3

REGION = 'ap-northeast-2'

def get_client(service_name='firehose', is_in_aws=True): 
    # 1. 로컬 환경 (is_in_aws=False)
    # 로컬에 설정된 aws configure 프로필이나 환경변수를 통해 인증
    if not is_in_aws:
        session = boto3.Session(region_name=REGION)
        return session.client(service_name)    
    
    # 2. AWS 내부 환경 (is_in_aws=True)
    # EC2, CloudShell 등 IAM Role이 부여된 환경에서는 키 입력 없이 자동 인증
    return boto3.client(service_name, region_name=REGION)

# 로컬에서 테스트할 경우 get_client('firehose', False) 로 호출
firehose = get_client() 

4. 데이터 전송 (put_record)

클라이언트 객체가 생성되었으면, 직렬화된 로그 데이터를 Firehose Delivery Stream으로 쏜다.

def send_log():
    response = firehose.put_record(
        DeliveryStreamName='de-ai-06-an2-kdf-log-to-s3',
        Record={
            # 주의: 데이터의 끝을 알리는 개행문자(\n)를 반드시 추가해야 함
            'Data': make_one_log() + "\n" 
        }
    )
    print(f'전송결과 : {response}') # HTTP 200이면 정상

  • 엔지니어링 포인트: S3에 적재된 파일을 나중에 읽어 들일 때, 각 로그 객체가 줄바꿈으로 구분되어 있어야 (JSON Lines 형태) 파싱 에러가 나지 않는다. 따라서 'Data': make_one_log() + "\n" 처리는 사소하지만 매우 중요한 부분이다.

5. Firehose의 핵심: 버퍼링 (Size & Interval)

Firehose가 받은 데이터를 즉시 S3에 넣지 않는다고 했다. 그렇다면 언제 넣을까?
AWS 콘솔에서 Firehose를 생성할 때 설정한 Buffer size(크기)Buffer interval(시간) 조건에 따른다.

  • Buffer Size (예: 1MiB): 전송된 데이터가 버퍼 메모리 내에 1 메가바이트만큼 쌓이면 S3로 Flush(전송)한다.
  • Buffer Interval (예: 60초): 데이터가 용량을 다 채우지 못했더라도, 첫 데이터가 들어온 지 60초가 경과하면 S3로 Flush한다.

"둘 중 하나라도 먼저 조건을 만족하면(Whichever happens first)" S3로 데이터가 압축되어(또는 원본 그대로) 업로드된다. 이를 통해 파일의 크기를 적절히 키워 S3에 저장함으로써 다운스트림(Downstream)에서의 데이터 처리 효율을 극대화할 수 있다.

profile
develop myself

0개의 댓글