전체 그래프
Spark

Structured Streaming

data-engineeringsparkstreamingkafkareal-time

상위: Spark

요약

Spark Structured Streaming은 마이크로 배치 기반의 실시간 스트림 처리 프레임워크입니다. Kafka와 통합하여 스트리밍 데이터를 처리하고, DataFrame API를 사용하여 배치 처리와 동일한 방식으로 코딩할 수 있습니다.

스트리밍 처리 프레임워크 비교

  • Flink → 진정한 실시간 / 초저지연 처리
    • 금융, 보안 등 즉시 반응이 중요한 분야에 적합
  • Spark → 기존 Spark 생태계 기반 + 실시간 분석도 가능
    • 데이터 분석, 추천 시스템, 로그 분석 등 폭넓은 응용에 적합

Flink vs Spark 상세 비교

항목Apache FlinkApache Spark Structured Streaming
지연 시간수 밀리초 수준 (초저지연)수십~수백 밀리초 (실시간성 양호)
처리 방식이벤트 기반 (True Streaming)마이크로 배치 기반 (Micro-batching)
정확성 보장기본 Exactly-once기본 At-least-once, 설정 시 Exactly-once 가능
API 지원Java, Scala, Table API, SQLPython(PySpark), Scala, SQL, R 등 풍부한 API
생태계 통합Flink 자체 생태계 중심Spark MLlib, GraphX, Delta Lake 등과 통합 용이
학습 난이도비교적 높은 진입 장벽PySpark 등으로 시작하기 쉬움
확장성실시간 처리에 최적화실시간 + 배치 통합 파이프라인 구성에 용이
대표 사용 예시Fraud detection, 실시간 로그 분석로그 집계, 클릭스트림 분석, ML 파이프라인 연계

Kafka + Spark 통합 실습

Kafka 메시지를 Spark로 읽고, 타임스탬프별 메시지 개수 집계 수행

실습 전 준비

Kafka topic 존재 확인

kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list

Spark 3.5.4 설치 및 PySpark 환경 설정

pip install pyspark==3.5.4
pip install findspark==2.0.1
  • 필요한 패키지는 spark.jars.packages 옵션으로 자동 설치됨

Kafka 메시지 전송 (spark_producer.py)

from kafka import KafkaProducer
import json
import random
import time
from datetime import datetime, timedelta

producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

base_time = datetime.now()

for i in range(30):
    offset_minutes = random.randint(0, 4)
    offset_seconds = random.randint(0, 59)

    msg_time = base_time + timedelta(minutes=offset_minutes, seconds=offset_seconds)

    message = {
        'id': i,
        'message': f'테스트 메시지 {i}',
        'timestamp': msg_time.strftime('%Y-%m-%d %H:%M:%S')
    }

    producer.send('test-topic', message)
    print(f'전송된 메시지: {message}')
    time.sleep(0.5)

producer.close()

Spark Kafka 집계 처리 (spark_kafka_example.py)

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window
from pyspark.sql.types import StructType, StructField, StringType, TimestampType

# SparkSession 생성 (Kafka 패키지 포함)
spark = SparkSession.builder \
    .appName("KafkaSparkIntegration") \
    .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.4") \
    .getOrCreate()

# Kafka에서 스트림 읽기
df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "test-topic") \
    .option("startingOffsets", "earliest") \
    .load()

def process_kafka_stream():
    # 메시지를 문자열로 변환
    value_df = df.selectExpr("CAST(value AS STRING)")

    # JSON 스키마 정의
    schema = StructType([
        StructField("timestamp", TimestampType(), True),
        StructField("message", StringType(), True)
    ])

    # JSON 파싱
    parsed_df = value_df.select(from_json(col("value"), schema).alias("data")).select("data.*")

    # 타임스탬프별 집계
    result_df = parsed_df \
        .withWatermark("timestamp", "1 minute") \
        .groupBy(window("timestamp", "1 minute")) \
        .count()

    # 콘솔에 출력
    query = result_df.writeStream \
        .outputMode("update") \
        .format("console") \
        .option("truncate", "false") \
        .start()

    query.awaitTermination()

if __name__ == "__main__":
    process_kafka_stream()

실습 순서

  1. 터미널 1: Kafka Producer 실행

    python spark_producer.py
  2. 터미널 2: Spark Streaming 작업 실행

    python spark_kafka_example.py
  3. 터미널 3: Kafka Consumer로 확인 (선택)

    /usr/local/kafka$ bin/kafka-console-consumer.sh \
      --topic test-topic \
      --bootstrap-server localhost:9092 \
      --from-beginning

결과 예시

Structured Streaming

Structured Streaming

+------------------------------------------+-----+
|window                                    |count|
+------------------------------------------+-----+
|{2024-10-14 10:00:00, 2024-10-14 10:01:00}|8    |
|{2024-10-14 10:01:00, 2024-10-14 10:02:00}|12   |
|{2024-10-14 10:02:00, 2024-10-14 10:03:00}|7    |
|{2024-10-14 10:03:00, 2024-10-14 10:04:00}|3    |
+------------------------------------------+-----+

Structured Streaming 핵심 개념

1. Watermark

늦게 도착하는 데이터를 처리하는 기준

df.withWatermark("timestamp", "10 minutes")
  • 10분 이전 데이터는 무시
  • 메모리 효율성 향상

2. Output Mode

  • append: 새로운 행만 출력
  • update: 변경된 행만 출력
  • complete: 전체 결과 출력
query = df.writeStream \
    .outputMode("update") \
    .format("console") \
    .start()

3. Trigger

실행 주기 설정

# 마이크로 배치 (기본)
.trigger(processingTime='5 seconds')

#  번만 실행
.trigger(once=True)

# 가능한 빠르게
.trigger(continuous='1 second')

4. Checkpoint

장애 복구를 위한 상태 저장

query = df.writeStream \
    .option("checkpointLocation", "/path/to/checkpoint") \
    .start()

다양한 Sink 타입

Console Sink (개발/테스트)

query = df.writeStream \
    .format("console") \
    .start()

File Sink (Parquet, CSV, JSON)

query = df.writeStream \
    .format("parquet") \
    .option("path", "/path/to/output") \
    .option("checkpointLocation", "/path/to/checkpoint") \
    .start()

Kafka Sink

query = df.writeStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("topic", "output-topic") \
    .start()

Memory Sink (테스트)

query = df.writeStream \
    .format("memory") \
    .queryName("my_table") \
    .start()

# 쿼리로 확인
spark.sql("SELECT * FROM my_table").show()

실전 패턴

1. 윈도우 기반 집계

# 5분 윈도우, 1분 슬라이드
windowed_counts = df \
    .groupBy(
        window("timestamp", "5 minutes", "1 minute"),
        "user_id"
    ) \
    .count()

2. Join (Stream-Static)

# 스트림과 정적 데이터 조인
enriched = streaming_df.join(static_df, "user_id")

3. Join (Stream-Stream)

#  스트림 조인
joined = stream1.join(
    stream2,
    expr("""
        stream1.user_id = stream2.user_id AND
        stream1.timestamp >= stream2.timestamp AND
        stream1.timestamp <= stream2.timestamp + interval 1 hour
    """)
)

4. Deduplication

# 중복 제거
deduplicated = df.dropDuplicates(["user_id", "event_id"])

모니터링

스트리밍 메트릭 확인

# 실행 중인 쿼리 상태
query.status

# 최근 진행 상황
query.recentProgress

# 마지막 진행 상황
query.lastProgress

Spark UI

  • http://localhost:4040
  • Streaming 탭에서 실시간 메트릭 확인
  • Input Rate, Process Rate, Batch Duration 등

성능 최적화

1. 적절한 배치 인터벌

# 처리 시간보다 약간  인터벌 설정
.trigger(processingTime='10 seconds')

2. 파티션 수 조정

spark.conf.set("spark.sql.shuffle.partitions", "200")

3. Watermark 설정

# 늦게 도착하는 데이터 처리
df.withWatermark("timestamp", "10 minutes")

4. Checkpoint 위치

# 빠른 스토리지 사용 (SSD, S3)
.option("checkpointLocation", "s3://bucket/checkpoint")