상위: Flink
요약
Flink는 실시간 스트리밍뿐만 아니라 배치 처리도 지원합니다. Table API를 사용하여 정적 CSV 데이터를 배치 처리하고, 카테고리별 통계를 계산할 수 있습니다. 스트리밍과 배치 모두 동일한 API로 처리 가능합니다.
Flink Batch Processing
- 실시간으로 처리하는 스트리밍 작업뿐 아니라 모아놓은 데이터로 실행하는 배치 작업도 처리 가능
- Table API를 사용하여 정적 CSV 데이터를 배치 처리
- 판매 데이터를 집계하여 카테고리별 통계를 계산
- 파일 기반 소스와 싱크를 통해 Flink 배치 파이프라인 흐름 학습
실습 의미 요약
| 항목 | 내용 |
|---|---|
| Flink 사용 방식 | 스트리밍 처리뿐만 아니라, 배치 처리도 동일한 Table API로 구현 가능 |
| Table API | SQL처럼 사용 가능하며, 대규모 데이터 처리에 최적화 |
| 입력 데이터 | CSV 포맷의 정적 판매 이력 데이터 (1000건 랜덤 생성) |
| 출력 데이터 | 카테고리별: 매출합, 평균가격, 수량합계, 거래수 계산 결과 |
| 사용된 기술 요소 | CREATE TABLE, INSERT INTO, GROUP BY, SUM, AVG 등의 SQL 연산 |
| 활용 시나리오 | 월별 판매 요약, 마케팅 리포트 생성, 데이터 마트 적재 등 |
실행 흐름
- 샘플 데이터 생성
- CSV 파일 저장 (source)
- Flink Table API로 로드
- SQL로 집계
- 결과를 싱크 디렉토리에 CSV로 저장
- 결과 파일 → 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'
관련 주제
- Flink - Flink 개요
- Kafka Integration - 실시간 처리
- Batch vs Real-time - 배치/실시간 비교