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

Idempotent Producer

Idempotent producer ติด tag ให้ทุก record ด้วย producer id, epoch และ sequence number ต่อ partition เพื่อให้ broker ทิ้ง duplicate ที่เกิดจาก retry ได้ — ให้การส่งแบบ exactly-once ต่อ partition

ถ้าไม่มี idempotence การ retry อาจสร้าง record ซ้ำเงียบ ๆ producer ส่ง record broker เขียนสำเร็จ แต่ acknowledge หายระหว่างทางกลับ producer timeout แล้ว retry และ broker ก็เขียน record เดิมเป็นครั้งที่สอง ตอนนี้เรามีสองสำเนา แถมแยกไม่ออกด้วย

idempotence แก้เรื่องนี้: broker assign producer id (PID) และ epoch ให้แต่ละ producer และ producer ก็ stamp ทุก record ด้วย sequence number ที่เพิ่มขึ้นเรื่อย ๆ ต่อ partition เมื่อ record มาถึง broker เช็ค sequence number: ถ้าเคยเห็นเลขนี้สำหรับ PID และ partition นั้นแล้ว จะ acknowledge แต่ไม่เขียน record ซ้ำอีก

# The Java client: idempotence is on by default; shown explicitly here
enable.idempotence=true
# acks=all is required and applied automatically once idempotence is on
acks=all
# retries must be > 0 (defaults high) so the producer actually retries
retries=2147483647
# ordering is preserved only if in-flight requests per connection stays <= 5
max.in.flight.requests.per.connection=5

KafkaJS ต้อง opt in เองที่ producer แล้ว client จะบังคับใช้ข้อกำหนดที่เทียบเท่ากันภายใน — คุณไม่ต้องตั้ง acks หรือ in-flight limit แยกเอง:

// idempotent: true implies acks: -1 (all) internally, and caps in-flight requests
const producer = kafka.producer({ idempotent: true, maxInFlightRequests: 5 })

idempotence ไม่ได้ยืนได้ด้วยตัวเอง แต่ต้องการ config เฉพาะเพื่อให้การรับประกันเป็นจริง:

  • acks=all — record ต้องถูกยืนยันโดยทุก in-sync replica
  • retries > 0 — retry คือต้นเหตุของ duplicate ตั้งแต่แรก ถ้าไม่มี retry ก็ไม่มีอะไรให้ dedupe
  • max.in.flight.requests.per.connection <= 5 — ยอมให้มี batch in-flight ได้ถึงห้าตัว โดยยังรักษา ordering ไว้ เพราะ broker ติดตาม sequence number และปฏิเสธอะไรที่มาผิดลำดับ

เงื่อนไขเหล่านี้เป็นข้อกำหนดระดับ protocol ไม่ใช่แค่เรื่องจุกจิกของ Java client เท่านั้น: ใน KafkaJS การส่ง idempotent: true ทำให้ client บังคับใช้การรับประกันที่เทียบเท่ากันภายใน นั่นคือเหตุผลที่คุณไม่ต้องตั้ง acks เองบน idempotent producer ของ KafkaJS

เมื่อเงื่อนไขเหล่านี้เป็นจริง เราจะได้ exactly-once delivery ต่อ partition ตอน retry: ไม่ว่า record จะถูก retry กี่ครั้ง ก็ถูกเขียนลง log ครั้งเดียว ตัว epoch คือสิ่งที่ทำให้ปลอดภัยข้ามการ restart ของ producer และ network partition — epoch ใหม่กว่าจะ fence producer instance เก่าที่อาจเป็น zombie ทำให้ write ค้างเก่าถูกปฏิเสธ

flowchart LR
  send1["Send seq=7 (PID 100, epoch 3)"] --> write["Broker writes seq=7"]
  write --> lostack["Ack lost on the way back"]
  lostack --> retry["Producer retries seq=7"]
  retry --> check["Broker: already saw seq=7 for this PID"]
  check --> dedupe["Acknowledge but do NOT write again"]
The broker dedupes a retried record by its sequence number

จุดนี้คนสับสนบ่อย idempotence รับประกันว่า record ถูกเขียน ครั้งเดียวต่อ partition แม้มี retry — แต่ทำได้ แค่นั้นไม่ได้ ทำให้การ write หลาย record หลาย partition เป็น atomic และ ไม่ได้ ผูก write ของ producer เข้ากับการ commit offset ของ consumer

การรับประกันที่แข็งกว่านั้น — การ write แบบ atomic ข้าม partition และ exactly-once processing จริง ๆ — คือสิ่งที่ transactions ให้ ซึ่งสร้าง ทับบน idempotent producer อีกที โมดูล Delivery Semantics และ Transactions จะครอบคลุม transactionalId, producer.transaction() และ txn.commit() ตอนนี้จำเส้นแบ่งไว้: idempotence กำจัด duplicate จาก retry ส่วน transactions เพิ่ม atomicity

idempotent producer ป้องกันอะไร
broker รู้จัก duplicate จาก retry ได้อย่างไร
idempotence เปิดเป็น default ใน KafkaJS ไหม
idempotence ต่างจาก transactions อย่างไร