โหมดมืด
บทที่ 11 — Stream Processing (Kafka Streams + Flink + ksqlDB)
← บทที่ 10: Broker Comparison | สารบัญ | บทที่ 12: Event-Driven Patterns →
TL;DR: บทนี้ตอบ "ประมวลผล event ต่อเนื่อง (ไม่ใช่ batch) ยังไง" — สอน concept ก่อน แล้วค่อยลง implement. เครื่องมือหลักที่จะเจอ: Kafka Streams (library ฝัง Java app), Apache Flink (cluster แยก), ksqlDB (SQL-on-streams). ข้ามได้ถ้า: ไม่ทำ real-time analytics / materialized view (มุมมองที่คำนวณไว้ล่วงหน้า พร้อม query) / fraud detection
Stream Processing คือการประมวลผล event ต่อเนื่อง ทันทีที่มาถึง — ต่างจาก batch processing ที่รวบข้อมูลเป็นก้อนแล้วประมวลทีหลังเป็นรอบ ๆ (เช่นรันทุกเที่ยงคืน)
📖 กางศัพท์:
- Stream Processing (การประมวลผลแบบลำดับต่อเนื่อง — event เข้า → process ทันที) ต่างจาก batch processing (รวมเป็นกอง แล้วประมวลรอบเดียวต่อชั่วโมง/วัน)
- event time (เวลาที่ event เกิดจริง — เช่น timestamp ตอน user คลิก) vs processing time (เวลาที่ stream processor "เห็น" event นั้น) — ทั้งคู่ต่างกันได้มากเมื่อมี network delay
- windowing (การแบ่งหน้าต่างเวลา — แบ่ง stream เป็นช่วง ๆ เพื่อ aggregate; tumbling = หน้าต่างไม่ทับกัน, sliding = เลื่อนทับกัน, session = แบ่งตามช่วงที่ user active)
- watermark (ลายน้ำ — เส้นเวลา event ที่ระบบมั่นใจแล้วว่า "ไม่มี event ก่อนเวลานี้จะมาเพิ่มอีก" → window ก่อนหน้าปิดได้)
- state store (คลังเก็บ state — local RocksDB ใน stream app) + changelog (สมุดบันทึก state changes ที่ replicate ไป Kafka topic เพื่อ recover state ตอน app restart)
- materialized view (มุมมองที่ pre-compute ค่าไว้พร้อม query — ต่างจาก view ปกติที่คำนวณ on-the-fly ทุกครั้ง)
ตัวอย่าง:
- Fraud detection real-time
- Live dashboard
- Recommendation update
- Materialized view ของ event sourcing (view ที่ pre-compute ค่าไว้พร้อม query)
- Anomaly detection
3 ตัวหลักที่บทนี้จะสอน:
- Kafka Streams — library ใน Java app (in-process = รันในกระบวนการเดียวกับ app ไม่ต้องตั้ง cluster แยก)
- Apache Flink — distributed framework, stateful, current standard สำหรับ large-scale
- ksqlDB — SQL-on-streams (Confluent project — เขียน
CREATE STREAM ... SELECT ...แบบ SQL แทน Java code, เบื้องหลังคือ Kafka Streams runtime)
(ยังมีตัวเลือกอื่นอีก เช่น Apache Beam, Materialize, RisingWave, Spark Structured Streaming, Arroyo — ดูรายละเอียดเปรียบเทียบทั้งหมดใน Part 9)
ใช้เวลา 3-4 ชั่วโมง
Part 1: Stream vs Batch
text
Batch:
Schedule (3 AM daily) → query DB ──→ result
Stream:
Event arrives → process incrementally → result update| Batch | Stream | |
|---|---|---|
| Latency | Hours | ms - sec |
| Reprocess | Re-run | Replay events |
| State | DB-stored | In-memory + checkpoints |
| Complexity | Low | Higher |
| Use case | Reports, ETL | Real-time analytics |
Part 2: Concepts
2.1 Event Time vs Processing Time
หัวใจของ stream processing คือต้องแยกให้ออกระหว่างเวลาสองแบบ: event time คือเวลาที่ event เกิดขึ้นจริง (เช่น timestamp ตอน user คลิก) ส่วน processing time คือเวลาที่ stream processor ได้รับ event นั้น ปกติสองเวลานี้ต่างกันเล็กน้อย แต่ network delay อาจทำให้ event มาถึงช้าหรือมาผิดลำดับได้มาก ระบบที่ดีจึงต้อง handle late event และ out-of-order event โดยยึด event time เป็นหลัก ไม่ใช่ processing time:
text
Event time: เวลาที่ event "เกิดจริง" (จาก timestamp ใน payload)
Processing time: เวลาที่ stream processor receiveEvent time > Processing time = late event (network delay)
→ Stream processor ต้อง handle late event + out-of-order
2.2 Window
เนื่องจาก stream ไม่มีจุดจบ การ aggregate ต้องแบ่งเวลาเป็น "หน้าต่าง" (window) — มีหลายแบบตามความต้องการ: tumbling (ช่วงไม่ทับกัน), hopping (เลื่อนทับกัน), sliding (ต่อ event), และ session (แบ่งตามช่วงเงียบ) เลือกให้ตรงกับคำถามที่อยากตอบ:
Tumbling Window (non-overlapping) — แบ่งหน้าต่างไม่ทับกัน เหมือนตัดเค้ก:
→ 1 event เข้า 1 window เท่านั้น เช่น "ยอด order ต่อ 5 นาที"
Hopping Window (overlap) — หน้าต่างเลื่อนทับกัน:
→ 1 event เข้าได้หลาย window (overlap) เช่น "moving average 5 นาที update ทุก 1 นาที"
Sliding Window (per event) — หน้าต่างเลื่อนตามแต่ละ event:
text
event @ t → window [t-5min, t]→ 1 event = 1 window ใหม่ มองย้อนหลังไป 5 นาที
Session Window (gap-based) — แบ่งตามช่วงเงียบ:
text
events spaced < 30 sec = same session
gap > 30 sec = new session→ เหมาะกับ "1 session ของ user คลิกอะไรบ้าง" — ระบบหา gap เอง
2.3 Watermark
Mark "เราเห็น event time ถึงตรงนี้ — late event หลังจากนี้ skip"
text
Events: [t=1, t=3, t=5, t=2 (late!), t=7, t=4 (very late!)]
Watermark = max event time - allowed lateness (e.g., 10 sec)
ถ้า watermark = 5, event t=2 ยังรับได้ (within lateness)
ถ้า watermark = 10, event t=4 = ต้อง drop2.4 State
Stream processor ต้องเก็บ state:
- Aggregation count (per key, per window)
- Join state (รอ match)
- Deduplication set
- Session state
→ State stored in:
- Memory (fast, ใหญ่ไม่ได้)
- RocksDB (embedded key-value database — ฐานข้อมูล key-value ที่ฝังในตัว app เอง เก็บลง disk บนเครื่องเดียวกัน ไม่ต้องตั้ง server แยก; Kafka Streams ใช้เป็น default) — disk-backed
- External (Flink + state backend)
→ ต้อง checkpoint เพื่อ recovery
Part 3: Kafka Streams
3.1 ทำไม Kafka Streams
- Java library — ไม่ต้อง cluster แยก (รันใน Spring Boot ได้)
- Built on Kafka — ใช้ Kafka เป็น state store backup
- Exactly-once semantics
- Scale by partition
3.2 Setup
เริ่มใช้ Kafka Streams — เพิ่ม dependency kafka-streams คู่กับ spring-kafka แล้วตั้ง application-id (ใช้ระบุ instance group + ตั้งชื่อ state store/internal topic) จุดสำคัญคือเปิด exactly_once_v2 เพื่อรับประกัน process ไม่ซ้ำ/ไม่หาย:
xml
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<!-- version managed by spring-boot-dependencies BOM -->
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<!-- version managed by spring-boot-dependencies BOM
ถ้าใช้นอก Spring Boot ต้องระบุ <version> ให้ match กับ broker -->
</dependency>yaml
spring.kafka:
streams:
application-id: order-aggregator
bootstrap-servers: kafka:9092
properties:
processing.guarantee: exactly_once_v2
num.stream.threads: 4⚠️ exactly_once_v2 requirements: ต้องการ broker version ≥ 2.5. Legacy
exactly_onceถูก deprecated ตั้งแต่ Kafka Streams 3.0 — ห้ามใช้ในงานใหม่ (จำ version broker ตัวเองไม่ได้? ย้อนไปดู timeline เวอร์ชันที่ บทที่ 08: Kafka Deep Dive)
3.3 Basic Topology
"topology" คือ pipeline ที่บอกว่าข้อมูลไหลจาก topic ต้นทาง ผ่านการแปลง (filter/map) ไปยัง topic ปลายทางอย่างไร — เขียนแบบ declarative ด้วย StreamsBuilder ตัวอย่างนี้อ่านจาก orders, กรองยอด > 100, enrich แล้วเขียนลง topic ใหม่ (สังเกต orderSerde()/enrichedSerde() — เป็น Serde (Serializer/Deserializer — ตัวแปลง object ↔ bytes เวลาส่งเข้า/ออก Kafka) ที่ต้องประกาศเป็น bean เอง เช่น new JsonSerde<>(OrderEvent.class) ไม่งั้น compile ไม่ผ่าน):
java
@Configuration
@EnableKafkaStreams
public class StreamConfig {
@Bean
public KStream<String, OrderEvent> orderStream(StreamsBuilder builder) {
KStream<String, OrderEvent> orders = builder.stream("orders",
Consumed.with(Serdes.String(), orderSerde()));
orders
.filter((k, v) -> v.amount() > 100)
.mapValues(v -> v.toEnriched())
.to("orders-enriched", Produced.with(Serdes.String(), enrichedSerde()));
return orders;
}
}3.4 Aggregation + Window
ตัวอย่างรวมแนวคิด window กับ aggregation เข้าด้วยกัน — นับ order ต่อลูกค้าในหน้าต่างเวลา 5 นาที โดยมี grace period เผื่อ late event ผลลัพธ์ออกเป็น KTable ที่อัปเดต real-time แล้วเขียนกลับเป็น stream (ตัวอย่างนี้ละ Consumed.with(...) ตอนอ่าน stream — ใช้ได้จริงเมื่อตั้ง default.key.serde / default.value.serde ไว้ใน config กลางแล้วเท่านั้น มิฉะนั้นต้องระบุ Serde ตรง ๆ เหมือนตัวอย่าง 3.3):
java
KStream<String, OrderEvent> orders = builder.stream("orders");
KTable<Windowed<String>, Long> ordersPerCustomer = orders
.groupBy((k, v) -> v.customerId())
.windowedBy(TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(5),
Duration.ofMinutes(1))) // grace period
.count(Materialized.as("orders-count")); // ตั้งชื่อ store ไว้ query ทีหลัง (ดู 3.7)
ordersPerCustomer.toStream()
.map((winKey, count) -> KeyValue.pair(
winKey.key(),
new OrderCountResult(winKey.key(), winKey.window().start(), count)))
.to("orders-per-customer-5min");3.5 Join Streams
stream-stream join คือการรวม 2 สตรีมที่สัมพันธ์กันด้วย key เดียวกันภายในกรอบเวลา (join window) — เช่นจับคู่ order กับ payment ที่เกิดใกล้กัน เพื่อสร้าง event ที่สมบูรณ์ขึ้น ต้องกำหนด window เพราะ event อาจมาไม่พร้อมกัน:
java
KStream<String, Order> orders = builder.stream("orders");
KStream<String, Payment> payments = builder.stream("payments");
orders.join(payments,
(order, payment) -> new EnrichedOrder(order, payment),
JoinWindows.ofTimeDifferenceAndGrace(
Duration.ofMinutes(5),
Duration.ofMinutes(1)))
.to("orders-enriched");3.6 KTable + Stream-Table Join
ขณะที่ KStream คือ "ลำดับเหตุการณ์" KTable คือ "สถานะล่าสุดต่อ key" (เก็บค่าใหม่ทับค่าเก่า) — การ join stream กับ table ใช้ enrich event ด้วยข้อมูลอ้างอิง เช่นเติมข้อมูลลูกค้าเข้าไปในทุก order โดยดึงจาก KTable ของ customer:
KTable = state ที่อัพเดตจาก stream:
text
order-1 → status=PLACED
order-1 → status=PAID
order-1 → status=SHIPPED
KTable:
order-1 = SHIPPED (latest)java
KTable<String, Customer> customers = builder.table("customers");
KStream<String, Order> orders = builder.stream("orders");
orders.leftJoin(customers,
(order, customer) -> new OrderWithCustomer(order, customer))
.to("orders-with-customer");3.7 Interactive Queries
จุดเด่นที่หลายคนไม่รู้: state store ของ Kafka Streams query ได้โดยตรงเหมือน DB — เปิด REST endpoint อ่านค่า aggregate ล่าสุดได้เลยโดยไม่ต้องเขียนผลลง DB ภายนอก ทำให้ stream processor ทำหน้าที่เป็น materialized view ที่อัปเดต real-time:
Read state store จาก external:
java
@RestController
class StateController {
@Autowired StreamsBuilderFactoryBean factoryBean;
@GetMapping("/count/{customerId}")
public Long getCount(@PathVariable String customerId) {
ReadOnlyKeyValueStore<String, Long> store = factoryBean.getKafkaStreams()
.store(StoreQueryParameters.fromNameAndType(
"orders-count", QueryableStoreTypes.keyValueStore()));
return store.get(customerId);
}
}→ stream processor ก็ทำเหมือน DB (ที่ update real-time)
Part 4: Apache Flink
4.1 ทำไม Flink
- True streaming (event-by-event, ไม่ใช่ micro-batch)
- Stateful ที่ขนาดใหญ่ (TB)
- Exactly-once native
- SQL + DataStream API
- Connector เยอะ (Kafka, JDBC, Elasticsearch, Hudi, Iceberg)
- Checkpoint + savepoint สำหรับ recovery + upgrade
ใช้สำหรับ:
- Real-time analytics ขนาดใหญ่
- Complex event processing (CEP)
- Multi-source ETL
4.2 Architecture
Flink รันแบบ cluster แยก — JobManager เป็นตัวประสานงาน (รับ job, จัด schedule) ส่วน TaskManager หลายตัวเป็น worker ที่รัน operator ขนานกันใน slot งานหนึ่งถูกแปลงเป็น dataflow graph แล้วกระจายไปทำขนานข้าม TaskManager:
text
[JobManager] ← coordinator
│
coordinates
│
▼
[TaskManager 1] [TaskManager 2] [TaskManager 3] ← worker
slots slots slotsJob = Dataflow graph of operators Operators run parallel ใน TaskManager slots
4.3 Hello Flink
ตัวอย่าง job Flink แบบครบวงจร — อ่านจาก Kafka (พร้อมตั้ง watermark สำหรับ event time), จัดกลุ่มตาม customer, aggregate ในหน้าต่าง 5 นาที แล้วเขียนผลกลับ Kafka สังเกตว่า DataStream API ให้ควบคุมทุกขั้นแบบ programmatic:
java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<OrderEvent> orders = env.fromSource(
KafkaSource.<OrderEvent>builder()
.setBootstrapServers("kafka:9092")
.setTopics("orders")
.setGroupId("flink-orders")
.setStartingOffsets(OffsetsInitializer.earliest())
.setValueOnlyDeserializer(new OrderDeserializer())
.build(),
WatermarkStrategy.<OrderEvent>forBoundedOutOfOrderness(Duration.ofMinutes(1))
.withTimestampAssigner((event, ts) -> event.timestamp()),
"Kafka Source"
);
DataStream<EnrichedOrder> enriched = orders
.keyBy(OrderEvent::customerId)
// Flink 1.18+: ใช้ Duration แทน Time (class Time ถูก deprecate)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
.aggregate(new RevenueAggregator());
enriched.sinkTo(KafkaSink.<EnrichedOrder>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(...)
.build());
env.execute("Order Revenue");4.4 Flink SQL
ไม่อยากเขียน Java? Flink SQL ให้เขียน stream processing เป็น SQL ล้วน — นิยาม Kafka topic เป็น "table" (พร้อม watermark) แล้ว query/aggregate ด้วย SQL ปกติ เหมาะกับ data analyst ที่ถนัด SQL มากกว่าโค้ด:
sql
CREATE TABLE orders (
order_id STRING,
customer_id STRING,
amount DOUBLE,
event_time TIMESTAMP_LTZ(3),
WATERMARK FOR event_time AS event_time - INTERVAL '1' MINUTE
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);
CREATE TABLE customer_summary (
customer_id STRING,
window_start TIMESTAMP(3),
revenue DOUBLE,
order_count BIGINT,
PRIMARY KEY (customer_id, window_start) NOT ENFORCED
) WITH (
'connector' = 'kafka',
'topic' = 'customer-summary',
...
);
INSERT INTO customer_summary
SELECT
customer_id,
TUMBLE_START(event_time, INTERVAL '5' MINUTES),
SUM(amount) as revenue,
COUNT(*) as order_count
FROM orders
GROUP BY
customer_id,
TUMBLE(event_time, INTERVAL '5' MINUTES);4.5 Checkpoints + Savepoints
เพราะ stream job รันต่อเนื่องและมี state Flink ต้องมีกลไก recovery — checkpoint คือ snapshot อัตโนมัติเป็นระยะ (กู้คืนเมื่อ crash), ส่วน savepoint คือ snapshot ที่สั่งเองสำหรับ upgrade/migrate job โดยไม่เสีย state:
java
env.enableCheckpointing(60_000); // every 1 min
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000);
env.getCheckpointConfig().setCheckpointStorage("s3://my-bucket/flink-checkpoints");- Checkpoint = periodic snapshot (auto)
- Savepoint = manual trigger (สำหรับ upgrade / migrate)
bash
flink savepoint <job-id> s3://my-bucket/savepoints/
flink run -s s3://my-bucket/savepoints/savepoint-... -d my-job.jar→ stop + restart job ที่ state เดิม
Part 5: ksqlDB
5.1 SQL on Kafka
ksqlDB ทำให้เขียน stream processing บน Kafka เป็น SQL ได้โดยไม่ต้องตั้ง cluster Flink — นิยาม STREAM/TABLE จาก topic แล้ว aggregate ด้วย SQL พร้อม EMIT CHANGES ที่ส่งผลลัพธ์ออกแบบ real-time ทุกครั้งที่มีข้อมูลใหม่:
sql
CREATE STREAM orders (
order_id VARCHAR KEY,
customer_id VARCHAR,
amount DOUBLE
) WITH (KAFKA_TOPIC='orders', VALUE_FORMAT='AVRO');
CREATE TABLE orders_per_customer AS
SELECT customer_id, COUNT(*) AS count, SUM(amount) AS revenue
FROM orders
WINDOW TUMBLING (SIZE 5 MINUTES)
GROUP BY customer_id
EMIT CHANGES;
-- query
SELECT * FROM orders_per_customer WHERE customer_id = 'cust-1' EMIT CHANGES;→ ไม่ต้องเขียน Java — query แบบ SQL streaming
5.2 เมื่อใช้ ksqlDB
✅ Quick prototyping ✅ Data team / analyst ใช้ SQL ได้ ✅ Materialized view real-time
❌ Complex logic (ใช้ Kafka Streams / Flink) ❌ Multi-source (ksqlDB จำกัด Kafka)
Part 6: Materialize + Streaming Database
Materialize = Postgres-compatible DB ที่ data มาจาก stream (เบื้องหลังคือ Differential Dataflow — incremental view maintenance ที่อัปเดต view เฉพาะส่วนที่เปลี่ยน)
sql
CREATE SOURCE orders FROM KAFKA BROKER 'kafka:9092' TOPIC 'orders' FORMAT AVRO;
CREATE MATERIALIZED VIEW revenue_per_customer AS
SELECT customer_id, SUM(amount) FROM orders GROUP BY customer_id;
-- ตอนนี้ query ปกติ SQL:
SELECT * FROM revenue_per_customer WHERE customer_id = 'cust-1';View ถูก update real-time จาก stream
⚠️ syntax
CREATE SOURCE ... FROM KAFKA BROKERด้านบนเป็นของ Materialize รุ่นเก่า (self-host) ที่เลิกใช้แล้ว — รุ่น cloud ปัจจุบันของ Materialize มี syntax ต่างออกไป (ตรวจสอบ syntax ล่าสุดก่อนใช้จริง) ถ้าอยากได้ตัวที่ self-host ได้ฟรีตอนนี้ ดู RisingWave (Apache 2.0) แทน
ใช้แทน "BI tool + stream processor + DB" tier
📌 2026 status: Materialize Inc. pivot ไปเป็น enterprise / cloud-only ตั้งแต่ปี 2023-2024 (ไม่มี community edition ที่ใช้ฟรีใน production แบบเดิม). ทางเลือก open-source ที่คล้ายกัน:
- RisingWave (Apache 2.0) — streaming database, PostgreSQL-compatible, มี cloud และ self-host ฟรี
- Apache Flink SQL — ถ้าทีมรับ Flink cluster ได้
Part 7: Common Patterns
7.1 Enrichment
pattern ที่พบบ่อยสุดคือ enrichment — เติมข้อมูลให้ event ที่บางเบา เช่น order มีแค่ customerId ก็ join กับ lookup table (KTable ของ customer) เพื่อเติมชื่อ/ระดับสมาชิก ทำให้ event ปลายทางใช้งานได้ครบโดยไม่ต้องไปดึงข้อมูลทีหลัง:
Join event กับ lookup table:
text
orders ──join──→ customers ──→ enriched-orders
(KTable)7.2 Aggregation
aggregation คือการสรุปหลาย event เป็นตัวเลขเดียวต่อหน้าต่างเวลา — เช่นนับ order ต่อนาที, หายอดรวมต่อชั่วโมง เป็น pattern พื้นฐานสำหรับทำ metric/dashboard แบบ real-time:
text
orders ──window--group--count──→ orders-per-min7.3 Anomaly Detection
ตรวจจับความผิดปกติแบบ real-time — คำนวณค่าสถิติเคลื่อนที่ (rolling average/stddev) จากสตรีม แล้วแจ้งเตือนเมื่อค่าใหม่หลุดกรอบปกติ (เช่นยอดเกิน 3 เท่าของ stddev) ใช้กับ fraud detection หรือ monitoring:
text
orders ──→ rolling-stats──→ alert if amount > 3*stddev7.4 Sessionization
sessionization คือการจับกลุ่ม event ที่เกิดต่อเนื่องกันเป็น "session" โดยใช้ session window (แบ่งเมื่อมีช่วงเงียบนานพอ) — เช่นรวม click ของผู้ใช้คนหนึ่งเป็น 1 session การเข้าชม เพื่อวิเคราะห์พฤติกรรม:
text
clickstream ──session-window──→ user-sessions7.5 Pattern Matching (CEP)
Complex Event Processing (CEP) คือการตรวจจับ "ลำดับเหตุการณ์" ที่มีความหมาย ไม่ใช่แค่ event เดี่ยว ๆ — เช่น login fail หลายครั้งติดกันแล้วสำเร็จภายใน 1 นาที = สัญญาณ credential stuffing Flink มี CEP API ให้นิยาม pattern แบบนี้ได้ตรง ๆ:
text
login.failed → login.failed → login.success
within 1 min → ALERT "credential stuffing"Flink CEP API:
java
Pattern<Event, ?> pattern = Pattern.<Event>begin("failed1")
.where(SimpleCondition.of(e -> e.type().equals("LOGIN_FAILED")))
.next("failed2")
.where(SimpleCondition.of(e -> e.type().equals("LOGIN_FAILED")))
.next("success")
.where(SimpleCondition.of(e -> e.type().equals("LOGIN_SUCCESS")))
.within(Duration.ofMinutes(1)); // Flink 1.18+: Duration แทน TimePart 8: Production Concerns
8.1 State Size
เรื่องที่กระทบ production มากสุดคือขนาด state — operator ที่ stateful เก็บข้อมูลต่อ key ถ้า key มีจำนวนมาก (cardinality สูง เช่น key เป็น userId) state จะบวมจนเกิน memory ต้องวางแผนใช้ RocksDB, tier ลง S3 หรือใส่ TTL ให้ลบของเก่า:
Stateful operator (aggregation) เก็บ state per key. ถ้า key cardinality สูง → state ใหญ่
→ ใช้ RocksDB backend, tier ไป S3, หรือ TTL
8.2 Backpressure
เช่นเดียวกับ broker เมื่อ producer ผลิตเร็วกว่าที่ processor ตามทัน buffer จะบวม — Flink/Kafka Streams ใช้กลไก pull-based ที่ชะลอการอ่านอัตโนมัติเมื่อ downstream ช้า บวกกับ auto-scale เพื่อกระจายงานเพิ่ม:
Producer เร็วกว่า processor → buffer โต
→ pull-based + checkpoint alignment + auto-scale
8.3 Upgrade
การ upgrade stream app ที่มี state ต้องระวังไม่ให้ state หาย — Kafka Streams กู้ state จาก changelog topic ได้ตอน rolling restart ส่วน Flink ใช้ savepoint (snapshot → stop → ขึ้นเวอร์ชันใหม่ → resume) ทั้งคู่ควรทดสอบใน staging ก่อนเสมอ:
Kafka Streams: rolling restart (state restore จาก changelog topic)
Flink: savepoint → stop → new version → resume from savepoint→ test ใน staging ก่อน
8.4 Schema Evolution
นอกจาก schema ของ message แล้ว state ที่เก็บไว้ก็มี schema เช่นกัน — พอ schema ของ state เปลี่ยน (เพิ่ม/ลด field) ต้อง migrate ของเดิมให้เข้ากันได้ ไม่งั้น restore ไม่ขึ้น Flink มี state schema evolution ส่วน Kafka Streams พึ่ง serde ที่ evolve ได้ (Avro):
State schema เปลี่ยน → ต้อง migrate (Flink: state schema evolution, Kafka Streams: ใช้ JsonSerde + Avro)
8.5 Monitoring
- Records in/out per second
- Processing latency
- Lag (Kafka offset)
- Checkpoint duration + success
- State size
Part 9: Choosing — Kafka Streams vs Flink vs ksqlDB
ตัวเลือกทั้งหมด (2026 stack):
- Kafka Streams — library ใน Java app
- Apache Flink — distributed framework, stateful, current standard สำหรับ large-scale
- ksqlDB — SQL-on-streams
- Apache Beam — unified batch + stream API (รัน job เดียวบน Flink/Spark/Dataflow ก็ได้)
- Materialize — streaming SQL (Differential Dataflow, enterprise/cloud-only หลังปี 2023-2024)
- RisingWave — streaming database, Apache 2.0 — ทางเลือก open-source คล้าย Materialize
- Spark Structured Streaming — micro-batch (เหมาะถ้ามี Spark batch อยู่แล้ว)
- Arroyo — เครื่องมือใหม่ Rust-based, real-time stream SQL
Quick guide
- Kafka Streams — embed ใน Spring Boot, small-medium scale, state < 100GB ต่อ instance (sweet spot — ใหญ่กว่านี้ tune RocksDB ได้แต่ใช้ Flink + external state backend ดีกว่า)
- Flink — large scale, multi-source, complex CEP, very large state, savepoint สำหรับ upgrade
- ksqlDB — SQL, ad-hoc analyst usage, materialized view
- Apache Beam — unified batch+stream, รัน job เดียวบนหลาย runner (Flink/Spark/Dataflow)
- RisingWave / Materialize — streaming database, ใช้ SQL สร้าง materialized view real-time
- Spark Structured Streaming — มี Spark batch workload อยู่แล้ว, micro-batch OK
- Arroyo — Rust-based, real-time, น้ำหนักเบา
Part 9.5: Production War Stories
🔥 Uber (2021): Flink state ระเบิด 5TB ใน 24 ชม.
Stream join 2 topic, ใช้ session window (no timeout config) → user session บางคน "ไม่จบ" → state สะสม → checkpoint นาน → cluster ฟื้นไม่ได้
Lesson:
- State TTL บังคับ (
KeyedStateTtlConfig.newBuilder(Duration.ofDays(7))) - Session window ต้อง gap clear (
SessionWindows.withGap(Duration.ofHours(1))) (Flink 1.18+: ใช้DurationแทนTime— classTimeถูก deprecate แล้ว) - Monitor state size + alert
- Use incremental checkpoint + RocksDB (Flink) สำหรับ state ใหญ่
🔥 บริษัท adtech (2022): Kafka Streams + EOS rebalance ทำ 10 นาที downtime
ใช้ EOS (Exactly-Once Semantics = รับประกันประมวลผลครบครั้งเดียว) + many partitions (500) → rebalance ต้อง coordinate ใหญ่ → 10 นาที STW (Stop-The-World = หยุดรับงานชั่วคราวระหว่าง rebalance)
Lesson:
- Migrate ไป CooperativeStickyAssignor (incremental rebalance) — ใช้ได้ใน production ตอนนี้
- ลด
num.stream.threadsต่อ instance (มี instance มากขึ้นแทน) - KIP-848 (Kafka Improvement Proposal #848 — Next-Gen Consumer Rebalance Protocol ที่ย้าย logic rebalance ไปทำฝั่ง server แทน client — ดูบท 08) จะแก้ปัญหานี้ในระยะยาว
- ⚠️ ใน Kafka 4.0 สถานะคือ Early Access เท่านั้น (ดูบท 08) — ไม่ใช่ production-default
- คาดว่าจะ GA ในเวอร์ชันต่อ ๆ ไป (4.1 / 4.2) — รอ stable release ก่อนใช้ใน production
🔥 บริษัท fintech (2023): late event ทำ aggregation ผิด
Watermark ตั้ง 1 min late tolerance → event บางตัวมา 5 min late → drop → revenue summary ผิด
Lesson:
- Watermark ต้อง match กับ "worst-case lateness" ของ source
- Use side output สำหรับ late event (เก็บไว้ analyze)
- Allowed lateness + accumulator pattern สำหรับ correction
Part 9.6: 📋 Cheat Sheet
Framework selection
text
Small scale + Spring + simple → Kafka Streams (library)
Large scale + complex CEP → Apache Flink
SQL-only + analyst use → ksqlDB / Materialize
Already on Spark → Spark Structured Streaming
Cloud managed → AWS Kinesis Analytics / GCP DataflowWindow cheat sheet
text
Tumbling: [0-5m] [5-10m] [10-15m] ← no overlap, contiguous
Hopping: [0-5m][2-7m][4-9m] ← overlapping, fixed size
Sliding: window per event (continuous)
Session: [user activity] gap=30min ← dynamic based on eventEvent-time + watermark mental model
text
Event-time: "เมื่อเกิดเหตุการณ์จริง" (จาก data)
Processing-time: "เมื่อ stream engine เห็น" (clock)
Ingestion-time: "เมื่อ Kafka รับ"
→ Production มัก use EVENT-TIME (correct) + WATERMARK (tolerance)
→ Watermark = "หลัง time X จะไม่มี event ก่อน time X-tolerance อีก"State management baseline
text
☐ State backend: RocksDB (Flink) / RocksDB + changelog (KStreams)
☐ Incremental checkpoint (large state)
☐ Checkpoint interval: 10s-5min ตาม state size
☐ State TTL: set max retention per key
☐ Savepoint ก่อน upgrade / config change
☐ Monitor state size growth
☐ Test recovery from checkpoint (chaos test)Common pitfalls
text
❌ No watermark / late tolerance → out-of-order = data loss
❌ Unbounded state → OOM, slow checkpoint
❌ Side effect inside operator (call DB) → at-least-once = duplicate
❌ Wall-clock time in business logic → wrong with replay
❌ Lambda / job upgrade without savepoint → state lost
❌ Mix event-time + processing-time → inconsistentPart 10: Checkpoint
- Event time vs Processing time?
- Window types 4 แบบ?
- Watermark คืออะไร?
- Stateful operator เก็บ state ที่ไหน?
- Kafka Streams vs Flink — เลือกอันไหน?
- KTable ต่าง KStream ยังไง?
- Interactive Queries ใน Kafka Streams ทำอะไร?
- Savepoint ของ Flink ใช้ทำอะไร?
- ksqlDB เหมาะกับ user แบบไหน?
- Materialize เป็น DB ที่ทำอะไร?
Part 11: สรุปบทนี้
- Stream processing = ประมวลผล event ต่อเนื่อง, low latency
- Event time + Watermark = handle out-of-order events
- Window: tumbling, hopping, sliding, session
- Kafka Streams = library — embed ใน Spring Boot
- Flink = framework — large scale + complex CEP
- ksqlDB = SQL on stream
- State = stored in RocksDB + Kafka changelog topic
- Checkpoint/Savepoint = recovery + upgrade
บทถัดไป — Event-Driven Patterns (Saga, Outbox, Inbox, CDC) ← สำคัญมาก