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

Delivery Semantics

delivery semantics มีสามแบบ — at-most-once, at-least-once และ exactly-once — จากที่นั่งของ consumer คุณเลือกระหว่างสองแบบแรกได้จาก จังหวะที่คุณ commit เป็นหลัก แล้วไปถึงแบบที่สามด้วย idempotent processing หรือ transaction.

  • At-most-once — แต่ละ record ถูกส่งศูนย์หรือหนึ่งครั้ง record อาจ หาย แต่ไม่มีทางซ้ำ ได้แบบนี้จากการ commit offset ก่อน process: ถ้า crash หลัง commit แต่ยังทำไม่เสร็จ record จะถูกข้ามตอน restart.
  • At-least-once — แต่ละ record ถูกส่งหนึ่งครั้งขึ้นไป record อาจ ซ้ำ แต่ไม่มีทางหาย ได้แบบนี้จากการ commit หลัง process: crash ก่อน commit จะทำให้อ่าน batch ซ้ำ นี่คือดีฟอลต์ที่ใช้จริงในทางปฏิบัติ
  • Exactly-once — แต่ละ record ส่งผลต่อผลลัพธ์เพียงครั้งเดียว ไม่มีหายไม่มีซ้ำ แบบนี้ต้องการมากกว่าแค่จังหวะ commit.
flowchart TB
  poll["poll batch"] --> choice{"commit before or after processing?"}
  choice -->|before| amo["at-most-once: crash then skip (loss)"]
  choice -->|after| alo["at-least-once: crash then reprocess (duplicates)"]
  alo --> idem["make processing idempotent to reach effectively exactly-once"]
Commit placement selects the first two semantics

at-least-once บวกกับ idempotent processing คือเส้นทางที่พบบ่อยที่สุด: ออกแบบ side effect ให้การ process record เดิมซ้ำสองครั้งให้ผลเหมือนทำครั้งเดียว กลยุทธ์ที่ใช้ได้จริง:

  • Upsert ด้วย natural key แทนการ insert ตรง ๆ เพื่อให้ record ซ้ำเขียนทับแทนที่จะเพิ่มซ้ำ
  • Dedup ด้วย record id ที่ producer ประทับมา แล้วทิ้ง id ที่เคย apply ไปแล้ว
  • ทำ operation ปลายทางให้ commutative/idempotent (เช่น “set status = shipped” ไม่ใช่ “increment count”).

เส้นทางที่แข็งแรงกว่าในระดับ framework คือ Kafka transaction — เขียน output record และ commit input offset พร้อมกันแบบ atomic — ซึ่ง module ถัดไปจะครอบคลุมเต็ม ๆ transaction ให้ exactly-once สำหรับ pipeline แบบ consume-transform-produce ที่อยู่ภายใน Kafka.

อะไรคือลักษณะของ at-most-once delivery?
โดยทั่วไปคุณทำ at-least-once ได้อย่างไร?
วิธีที่ใช้ได้จริงในการทำให้ at-least-once ทำงานเหมือน exactly-once คืออะไร?
เส้นทาง exactly-once ที่แข็งแรงที่สุดระดับ framework ใน Kafka ใช้อะไร?