Skip to content

บทที่ 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
  • ksqlDBSQL-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
BatchStream
LatencyHoursms - sec
ReprocessRe-runReplay events
StateDB-storedIn-memory + checkpoints
ComplexityLowHigher
Use caseReports, ETLReal-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 receive

Event 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 = ต้อง drop

2.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

  • 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             slots

Job = Dataflow graph of operators Operators run parallel ใน TaskManager slots

ตัวอย่าง 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");

ไม่อยากเขียน 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-min

7.3 Anomaly Detection

ตรวจจับความผิดปกติแบบ real-time — คำนวณค่าสถิติเคลื่อนที่ (rolling average/stddev) จากสตรีม แล้วแจ้งเตือนเมื่อค่าใหม่หลุดกรอบปกติ (เช่นยอดเกิน 3 เท่าของ stddev) ใช้กับ fraud detection หรือ monitoring:

text
orders ──→ rolling-stats──→ alert if amount > 3*stddev

7.4 Sessionization

sessionization คือการจับกลุ่ม event ที่เกิดต่อเนื่องกันเป็น "session" โดยใช้ session window (แบ่งเมื่อมีช่วงเงียบนานพอ) — เช่นรวม click ของผู้ใช้คนหนึ่งเป็น 1 session การเข้าชม เพื่อวิเคราะห์พฤติกรรม:

text
clickstream ──session-window──→ user-sessions

7.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 แทน Time

Part 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
  • ksqlDBSQL-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
  • ksqlDBSQL, 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

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 — class Time ถูก 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 Dataflow

Window 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 event

Event-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       → inconsistent

Part 10: Checkpoint

  1. Event time vs Processing time?
  2. Window types 4 แบบ?
  3. Watermark คืออะไร?
  4. Stateful operator เก็บ state ที่ไหน?
  5. Kafka Streams vs Flink — เลือกอันไหน?
  6. KTable ต่าง KStream ยังไง?
  7. Interactive Queries ใน Kafka Streams ทำอะไร?
  8. Savepoint ของ Flink ใช้ทำอะไร?
  9. ksqlDB เหมาะกับ user แบบไหน?
  10. 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) ← สำคัญมาก


← บทที่ 10 | สารบัญ | บทที่ 12: Event-Driven Patterns →