전체 그래프
Flink

Batch Processing

data-engineeringflinkbatch-processingtable-api

상위: Flink

요약

Flink는 실시간 스트리밍뿐만 아니라 배치 처리도 지원합니다. Table API를 사용하여 정적 CSV 데이터를 배치 처리하고, 카테고리별 통계를 계산할 수 있습니다. 스트리밍과 배치 모두 동일한 API로 처리 가능합니다.

Flink Batch Processing

  • 실시간으로 처리하는 스트리밍 작업뿐 아니라 모아놓은 데이터로 실행하는 배치 작업도 처리 가능
  • Table API를 사용하여 정적 CSV 데이터를 배치 처리
  • 판매 데이터를 집계하여 카테고리별 통계를 계산
  • 파일 기반 소스와 싱크를 통해 Flink 배치 파이프라인 흐름 학습

실습 의미 요약

항목내용
Flink 사용 방식스트리밍 처리뿐만 아니라, 배치 처리도 동일한 Table API로 구현 가능
Table APISQL처럼 사용 가능하며, 대규모 데이터 처리에 최적화
입력 데이터CSV 포맷의 정적 판매 이력 데이터 (1000건 랜덤 생성)
출력 데이터카테고리별: 매출합, 평균가격, 수량합계, 거래수 계산 결과
사용된 기술 요소CREATE TABLE, INSERT INTO, GROUP BY, SUM, AVG 등의 SQL 연산
활용 시나리오월별 판매 요약, 마케팅 리포트 생성, 데이터 마트 적재 등

실행 흐름

  1. 샘플 데이터 생성
  2. CSV 파일 저장 (source)
  3. Flink Table API로 로드
  4. SQL로 집계
  5. 결과를 싱크 디렉토리에 CSV로 저장
  6. 결과 파일 → Pandas로 확인 (선택)

실습 코드 (pyflink_batch_example.py)

환경 설정

import os
import random
import logging
from datetime import datetime, timedelta
from pyflink.table import EnvironmentSettings, TableEnvironment

# 로깅 설정
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

def main():
    # 배치 실행 환경 설정
    env_settings = EnvironmentSettings.new_instance().in_batch_mode().build()
    table_env = TableEnvironment.create(env_settings)

샘플 데이터 생성

    # 데이터 디렉토리 생성
    data_dir = "flink_batch_data"
    output_dir = "flink_batch_output"
    os.makedirs(data_dir, exist_ok=True)
    os.makedirs(output_dir, exist_ok=True)
    
    # 샘플 데이터 생성
    categories = ["electronics", "peripherals", "accessories"]
    products = [f"product_{i}" for i in range(1, 51)]
    
    sample_data = []
    for _ in range(1000):
        category = random.choice(categories)
        product = random.choice(products)
        price = round(random.uniform(10.0, 1000.0), 2)
        quantity = random.randint(1, 5)
        date = datetime.now() - timedelta(days=random.randint(0, 30))
        
        sample_data.append(f"{category},{product},{price},{quantity},{date.strftime('%Y-%m-%d')}")
    
    # CSV 파일 저장
    csv_path = os.path.join(data_dir, "sales.csv")
    with open(csv_path, "w") as f:
        f.write("category,product,price,quantity,sale_date\n")
        f.write("\n".join(sample_data))
    
    logger.info(f"샘플 데이터 생성 완료: {csv_path}")

소스 테이블 정의

    # 소스 테이블 정의
    table_env.execute_sql(f"""
        CREATE TABLE sales (
            category STRING,
            product STRING,
            price DOUBLE,
            quantity INT,
            sale_date STRING
        ) WITH (
            'connector' = 'filesystem',
            'path' = '{data_dir}',
            'format' = 'csv',
            'csv.ignore-parse-errors' = 'true'
        )
    """)
    logger.info("소스 테이블 생성 완료")

싱크 테이블 정의

    # 싱크 테이블 정의
    table_env.execute_sql(f"""
        CREATE TABLE sales_summary (
            category STRING,
            total_revenue DOUBLE,
            avg_price DOUBLE,
            total_quantity BIGINT,
            transaction_count BIGINT
        ) WITH (
            'connector' = 'filesystem',
            'path' = '{output_dir}',
            'format' = 'csv'
        )
    """)
    logger.info("싱크 테이블 생성 완료")

SQL 집계 실행

    # SQL 실행
    logger.info("데이터 처리 쿼리 실행 시작")
    result = table_env.execute_sql("""
        INSERT INTO sales_summary
        SELECT
            category,
            SUM(price * quantity) AS total_revenue,
            AVG(price) AS avg_price,
            SUM(quantity) AS total_quantity,
            COUNT(*) AS transaction_count
        FROM sales
        GROUP BY category
    """)
    logger.info(f"쿼리 실행 결과: {result.get_job_client().get_job_status()}")
    print(f"배치 처리가 완료되었습니다. 결과는 {output_dir} 에 저장되었습니다.")

결과 확인

    # 결과 파일 읽기 (Pandas 사용)
    import pandas as pd
    import glob
    
    result_files = glob.glob(f"{output_dir}/*.csv")
    if result_files:
        df = pd.read_csv(result_files[0], header=None)
        df.columns = ["category", "total_revenue", "avg_price", "total_quantity", "transaction_count"]
        print("\n집계 결과:")
        print(df)

결과 출력 예시

category     total_revenue     avg_price     total_quantity     transaction_count
electronics     810835.53         486.07           1661                   338
peripherals     803306.04         499.90           1601                   322
accessories     830635.74         505.19           1673                   340

Flink 배치 처리 특징

1. 스트리밍과 동일한 API

# 배치 모드
env_settings = EnvironmentSettings.new_instance().in_batch_mode().build()

# 스트리밍 모드 (동일한 SQL!)
env_settings = EnvironmentSettings.new_instance().in_streaming_mode().build()

2. SQL 기반 처리

  • 복잡한 집계도 SQL로 간단히 표현
  • GROUP BY, JOIN, 윈도우 함수 모두 지원

3. 파일 기반 I/O

# 다양한 포맷 지원
'format' = 'csv'      # CSV
'format' = 'json'     # JSON
'format' = 'parquet'  # Parquet
'format' = 'avro'     # Avro

활용 시나리오

월별 판매 요약

SELECT 
    SUBSTRING(sale_date, 1, 7) as month,
    category,
    SUM(price * quantity) as monthly_revenue
FROM sales
GROUP BY SUBSTRING(sale_date, 1, 7), category

마케팅 리포트 생성

SELECT 
    category,
    COUNT(DISTINCT product) as product_count,
    AVG(price) as avg_price,
    SUM(quantity) as total_sold
FROM sales
WHERE sale_date >= '2024-10-01'
GROUP BY category

데이터 마트 적재

INSERT INTO daily_summary
SELECT 
    sale_date,
    category,
    SUM(price * quantity) as daily_revenue,
    COUNT(*) as transaction_count
FROM sales
GROUP BY sale_date, category

배치 vs 스트리밍

배치 처리 적합

  • 과거 데이터 재처리
  • 일별/월별 집계
  • 대용량 ETL

스트리밍 처리 적합

  • 실시간 대시보드
  • 즉각적인 알림
  • 연속적인 집계

통합 접근

# 동일한 쿼리로 배치/스트리밍 전환
def create_aggregation_query():
    return """
        SELECT category, COUNT(*) as cnt
        FROM source
        GROUP BY category
    """

# 배치 환경
batch_env = EnvironmentSettings.new_instance().in_batch_mode().build()
# 스트리밍 환경
stream_env = EnvironmentSettings.new_instance().in_streaming_mode().build()

성능 최적화

1. 파티셔닝

# 파티션 필드 기준으로 셔플  기록 (또는 CREATE TABLE의 PARTITIONED BY (category) 사용)
'sink.shuffle-by-partition.enable' = 'true'

2. 병렬 처리

# Parallelism 설정
table_env.get_config().get_configuration().set_string("parallelism.default", "4")

3. 메모리 관리

# 메모리 설정
'table.exec.resource.default-parallelism' = '8'

관련 주제