Spark Structured Streaming — Stream Processing

Spark Structured Streaming

Spark Structured Streaming Real-time Stream Processing Kafka Window Watermark State Management Production Deployment
| Feature | Structured Streaming | Kafka Streams | Flink |
|---|---|---|---|
| Processing | Micro-batch / Continuous | Record-by-record | Record-by-record |
| Latency | 100ms - seconds | ms | ms |
| Exactly-once | Yes (with checkpoint) | Yes | Yes |
| State | Built-in + RocksDB | Built-in + RocksDB | Built-in + RocksDB |
| SQL Support | Full SQL | KSQL (separate) | Flink SQL |
| Best For | Batch + Stream unified | Kafka-native apps | Low-latency stream |
เคล็ดลับ
- Watermark: ตั้ง Watermark เสมอ ป้องกัน State โตไม่จำกัด
- Checkpoint: ใช้ S3/HDFS สำหรับ Checkpoint ไม่ใช้ Local Disk
- RocksDB: ใช้ RocksDB State Store สำหรับ State ขนาดใหญ่
- Monitor: ดู Processing Time vs Trigger Interval ถ้าเกิน = Lag
- Delta: ใช้ Delta Lake เป็น Sink สำหรับ ACID + Time Travel
การนำความรู้ไปประยุกต์ใช้งานจริง

แหล่งเรียนรู้ที่แนะนำ ได้แก่ Official Documentation ที่อัพเดทล่าสุดเสมอ Online Course จาก Coursera Udemy edX ช่อง YouTube คุณภาพทั้งไทยและอังกฤษ และ Community อย่าง Discord Reddit Stack Overflow ที่ช่วยแลกเปลี่ยนประสบการณ์กับนักพัฒนาทั่วโลก
Structured Streaming คืออะไร
Spark SQL Stream Processing Micro-batch Continuous Exactly-once Kafka Delta Lake DataFrame API Catalyst Optimizer Event-time
Window Functions ใช้อย่างไร
Tumbling 5 นาทีไม่ซ้อน Sliding 10 นาที Slide 5 นาทีซ้อน Session Gap 10 นาที groupBy window session_window Metrics Average
Watermark ทำงานอย่างไร
Late Data จัดการข้อมูลสาย withWatermark timestamp duration State ไม่โต Memory ตัดข้อมูลเกิน Window ค่ามากใช้ Memory น้อยตัดข้อมูล
Production Deployment ทำอย่างไร
Checkpoint S3 HDFS Trigger Interval Kafka Source RocksDB State Monitor Spark UI Prometheus Alert Lag Delta Lake ACID Scale K8s YARN
สรุป
Spark Structured Streaming Kafka Window Tumbling Sliding Session Watermark State Checkpoint RocksDB Delta Lake Production Monitor





