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

การเลือก Processing Layer

ไม่มีวิธี “ที่ถูกต้อง” เพียงหนึ่งเดียวในการประมวลผล topic ของ Kafka — เลือก เครื่องมือที่เบาที่สุดที่พอดีกับงาน ตั้งแต่ plain consumer สำหรับการควบคุมเต็มที่ ไปจนถึง Kafka Streams สำหรับ logic แบบฝัง ksqlDB สำหรับ SQL แบบ declarative และ Connect สำหรับการย้ายข้อมูลล้วน ๆ

แต่ละ layer แลกความควบคุมกับความสะดวก:

  • plain consumer — คุณเขียน loop เองด้วย client ดิบ ควบคุมได้สูงสุด ใช้ภาษาไหนก็ได้ แต่คุณต้องดูแล offset, retry, state และ scaling เอง
  • Kafka Streamslibrary บน JVM ที่คุณฝังไว้ใน service เขียนโค้ดจริง (map, join, windowing) โดยเรื่อง state และ exactly-once ถูกจัดการให้ แต่ใช้ได้แค่ Java/Scala และอยู่ในแอปของคุณ
  • ksqlDB / SQL — เขียน stream processing เป็น SQL แบบ declarative (CREATE STREAM ... SELECT ...) เขียนได้เร็วสุด ไม่มีโค้ดให้ deploy แต่ถูกจำกัดด้วยสิ่งที่ SQL บอกได้
  • Kafka Connectย้ายข้อมูล ไม่มี logic ingest database หรือ sink ลง warehouse ด้วย config ล้วน ๆ
  • external stream processor — Flink, Spark Structured Streaming และเพื่อน ๆ cluster แยกสำหรับงานขนาดใหญ่มาก, semantics แบบ event-time ที่ซับซ้อน หรือเมื่อคุณต้องการ engine เดียวที่ครอบหลาย source
flowchart TD
  start["What does the job need?"] --> move{"Just move data,
no logic?"}
  move -->|"yes"| connect["Kafka Connect"]
  move -->|"no"| logic{"Express it in SQL?"}
  logic -->|"yes"| ksql["ksqlDB"]
  logic -->|"no"| jvm{"JVM app, embedded?"}
  jvm -->|"yes"| streams["Kafka Streams"]
  jvm -->|"no"| scale{"Huge scale or
non-JVM control?"}
  scale -->|"external engine"| flink["Flink / Spark"]
  scale -->|"full control"| consumer["Plain consumer"]
Pick by what the job needs

ไล่จากตัวที่พอดีและง่ายที่สุดลงมา:

  1. แค่ย้ายข้อมูล? ใช้ Connect — อย่าเขียนโค้ดเพื่อก๊อป table
  2. logic ที่เขียนเป็น SQL ได้? ใช้ ksqlDB — filter, join และ windowed aggregate โดยไม่ต้อง deploy
  3. logic ที่ต้องเขียนโค้ดจริง บน JVM? ใช้ Kafka Streams — ฝังในแอป, stateful, exactly-once และ scale ไปพร้อมแอป
  4. ต้องใช้ภาษาอื่น หรือต้องคุม loop เต็มที่? ใช้ plain consumer
  5. สเกลมหาศาลหรือการประมวลผล event-time ขั้นสูงข้ามหลายระบบ? หยิบ external processor อย่าง Flink มาใช้

ระบบที่ดีมักผสมกันหลายตัว: Connect ทำ ingest, Streams หรือ ksqlDB ทำ transform และ plain consumer ส่งช่วงสุดท้ายเข้า service กฎง่าย ๆ คือ: อย่าจ่ายค่าพลังที่คุณจะไม่ได้ใช้

ตอนนี้คุณมีชุดเครื่องมือประมวลผลครบแล้ว Schema Registry ให้ record ของคุณมี contract แบบมีเวอร์ชัน ทีมจึงวิวัฒน์ได้อย่างปลอดภัย Connect ย้ายข้อมูลเข้าออกโดยไม่ต้องเขียนโค้ด Kafka Streams ประมวลผล stream ด้วย logic แบบ stateful จริง ๆ ในแอปของคุณ และคู่มือนี้บอกคุณว่าควรหยิบ ตัวไหน มาใช้ เมื่อรวมกับโมดูลก่อนหน้าเรื่อง broker, producer, consumer และการรับประกันการส่ง คุณก็ออกแบบระบบ Kafka แบบ end-to-end ได้แล้ว — และอธิบายเหตุผลของทุก layer ในนั้นได้ด้วย

ทีมหนึ่งต้องก๊อป table ของ Postgres เข้า topic ของ Kafka โดยไม่มีการ transform เครื่องมือที่ดีที่สุดคือ
ksqlDB ให้คุณทำอะไรได้
ทำไมถึงเลือก plain consumer แทน Kafka Streams
กฎง่าย ๆ โดยรวมสำหรับการเลือก processing layer คืออะไร