ข้ามไปยังเนื้อหา

Kafka Streams

Kafka Streams คือ library ภาษา Java — ไม่ใช่ cluster แยกต่างหาก — ที่คุณฝังไว้ในแอปของตัวเองเพื่อ transform, aggregate และ join topic ของ Kafka แบบ stream ต่อเนื่อง โดยที่เรื่อง state และการรับประกัน exactly-once ถูกจัดการให้เรียบร้อย

Kafka Streams ให้มุมมองต่อ topic สองแบบ:

  • KStream — ลำดับ event ที่อิสระต่อกันและไม่มีจุดสิ้นสุด ทุก record คือ fact ใหม่ (“ผู้ใช้คลิก”) ไม่มีอะไรถูกเขียนทับ
  • KTablechangelog ที่ตีความว่าเป็นค่าล่าสุด ต่อ key แต่ละ record อัปเดต row ของ key นั้น เหมือนทำ UPSERT ลง table

นี่คือ stream-table duality: table คือ snapshot ที่ได้จากการเล่นซ้ำ stream ของ update ส่วน stream คือลำดับของการเปลี่ยนแปลงที่ได้จากการเฝ้าดู table topic เดียวกันอ่านได้ทั้งสองแบบขึ้นกับสิ่งที่คุณต้องการ — feed ของ event ที่ไหลอยู่ หรือ state ปัจจุบันต่อ key

flowchart LR
  s["KStream (events)"] -->|"aggregate / reduce"| t["KTable (latest per key)"]
  t -->|"toStream()"| s2["KStream (changes)"]
A stream of updates folds into a table; a table emits a stream of changes

operation แบ่งเป็นสองตระกูล:

  • stateless — จัดการแต่ละ record แยกกัน: map, filter, flatMap, branch ไม่ต้องจำ record ก่อนหน้า
  • stateful — ผลลัพธ์ขึ้นกับ record ที่เห็นมาแล้ว: aggregate, count, reduce, join และ windowing (จัดกลุ่ม record เป็นถังตามเวลา) พวกนี้ต้องมีที่เก็บผลลัพธ์ที่กำลังสะสม

ที่เก็บนั้นคือ state store — key-value store แบบ local (default คือ RocksDB) ที่อยู่ใน instance แอปของคุณ เพื่อให้รอดจากการ crash ทุก state store ถูกหนุนด้วย changelog topic ใน Kafka: ทุก update ถูกเขียนลง changelog ด้วย ดังนั้นถ้า instance restart ที่อื่น Streams จะสร้าง store ขึ้นใหม่ด้วยการเล่นซ้ำ changelog นั้น state ใน local ของคุณจึงทนทานโดยที่คุณไม่ต้องดูแล database เอง

นี่คือตัวอย่างคลาสสิก — นับว่าแต่ละคำปรากฏบ่อยแค่ไหนใน input topic:

StreamsBuilder builder = new StreamsBuilder();
// Read the input topic as a stream of (key, line) events
KStream<String, String> lines = builder.stream("text-input");
KTable<String, Long> counts = lines
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split("\\W+"))) // split into words
.groupBy((key, word) -> word) // re-key by the word itself
.count(); // stateful: running total per word -> a KTable
// Write the changelog of counts out to an output topic
counts.toStream().to("word-counts");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

count() สร้าง KTable ที่หนุนด้วย state store และ changelog ของตัวเอง ส่วน toStream().to(...) ส่งค่านับที่อัปเดตแล้วทุกครั้งต่อไปยัง downstream

stream processing ที่อ่าน อัปเดต state แล้วเขียน เป็นเรื่องที่ทำ ผิด ได้ง่ายเมื่อเกิด failure — คุณอาจนับซ้ำ Kafka Streams ย่อเรื่องนี้ให้เหลือแค่การตั้งค่าบรรทัดเดียว เมื่อตั้ง processing.guarantee เป็น exactly_once_v2 Streams จะห่อรอบ read-process-write แต่ละรอบ (รวมถึง changelog ของ state store และ consumer offset) ไว้ใน transaction เดียวของ Kafka การ retry จึงไม่ apply update ซ้ำสองครั้ง:

// Turn on end-to-end exactly-once for the whole topology
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, "exactly_once_v2");
// Modern Kafka: use the dedicated broker-side rebalance protocol for Streams groups
props.put(StreamsConfig.GROUP_PROTOCOL_CONFIG, "streams");
KStream กับ KTable ต่างกันอย่างไร
state store ของ Kafka Streams รอดจากการ crash ของแอปได้อย่างไร
ข้อใดเป็น operation แบบ STATEFUL
processing.guarantee=exactly_once_v2 ให้อะไร