상위: Spark
요약
Spark Structured Streaming은 마이크로 배치 기반의 실시간 스트림 처리 프레임워크입니다. Kafka와 통합하여 스트리밍 데이터를 처리하고, DataFrame API를 사용하여 배치 처리와 동일한 방식으로 코딩할 수 있습니다.
스트리밍 처리 프레임워크 비교
- Flink → 진정한 실시간 / 초저지연 처리
- 금융, 보안 등 즉시 반응이 중요한 분야에 적합
- Spark → 기존 Spark 생태계 기반 + 실시간 분석도 가능
- 데이터 분석, 추천 시스템, 로그 분석 등 폭넓은 응용에 적합
Flink vs Spark 상세 비교
| 항목 | Apache Flink | Apache Spark Structured Streaming |
|---|---|---|
| 지연 시간 | 수 밀리초 수준 (초저지연) | 수십~수백 밀리초 (실시간성 양호) |
| 처리 방식 | 이벤트 기반 (True Streaming) | 마이크로 배치 기반 (Micro-batching) |
| 정확성 보장 | 기본 Exactly-once | 기본 At-least-once, 설정 시 Exactly-once 가능 |
| API 지원 | Java, Scala, Table API, SQL | Python(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: Kafka Producer 실행
python spark_producer.py -
터미널 2: Spark Streaming 작업 실행
python spark_kafka_example.py -
터미널 3: Kafka Consumer로 확인 (선택)
/usr/local/kafka$ bin/kafka-console-consumer.sh \ --topic test-topic \ --bootstrap-server localhost:9092 \ --from-beginning
결과 예시


+------------------------------------------+-----+
|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")