Kafka Streams
ไอเดียในหนึ่งประโยค
หัวข้อที่มีชื่อว่า “ไอเดียในหนึ่งประโยค”Kafka Streams คือ library ภาษา Java — ไม่ใช่ cluster แยกต่างหาก — ที่คุณฝังไว้ในแอปของตัวเองเพื่อ transform, aggregate และ join topic ของ Kafka แบบ stream ต่อเนื่อง โดยที่เรื่อง state และการรับประกัน exactly-once ถูกจัดการให้เรียบร้อย
KStream, KTable และ duality
หัวข้อที่มีชื่อว่า “KStream, KTable และ duality”Kafka Streams ให้มุมมองต่อ topic สองแบบ:
KStream— ลำดับ event ที่อิสระต่อกันและไม่มีจุดสิ้นสุด ทุก record คือ fact ใหม่ (“ผู้ใช้คลิก”) ไม่มีอะไรถูกเขียนทับKTable— changelog ที่ตีความว่าเป็นค่าล่าสุด ต่อ 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)"]
operation แบบ stateless กับ stateful
หัวข้อที่มีชื่อว่า “operation แบบ stateless กับ stateful”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 เอง
topology นับคำ
หัวข้อที่มีชื่อว่า “topology นับคำ”นี่คือตัวอย่างคลาสสิก — นับว่าแต่ละคำปรากฏบ่อยแค่ไหนใน input topic:
StreamsBuilder builder = new StreamsBuilder();
// Read the input topic as a stream of (key, line) eventsKStream<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 topiccounts.toStream().to("word-counts");
KafkaStreams streams = new KafkaStreams(builder.build(), props);streams.start();count() สร้าง KTable ที่หนุนด้วย state store และ changelog ของตัวเอง ส่วน toStream().to(...) ส่งค่านับที่อัปเดตแล้วทุกครั้งต่อไปยัง downstream
exactly-once และ streams protocol
หัวข้อที่มีชื่อว่า “exactly-once และ streams protocol”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 topologyprops.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, "exactly_once_v2");
// Modern Kafka: use the dedicated broker-side rebalance protocol for Streams groupsprops.put(StreamsConfig.GROUP_PROTOCOL_CONFIG, "streams");