1 điểm bởi GN⁺ 2024-11-14 | 1 bình luận | Chia sẻ qua WhatsApp
  • Trong quá trình kiểm chứng Bufstream 0.1.0~0.1.3, một hệ thống streaming tương thích Kafka, đã phát hiện 2 vấn đề về tính sẵn sàng và 3 vấn đề về tính an toàn của chính Bufstream; cả 5 vấn đề đều đã được sửa tính đến bản 0.1.3
  • Bài kiểm thử dựa trên Java Kafka Client 3.8.0 và các bài kiểm thử Jepsen hiện có cho Kafka/Redpanda, đồng thời sử dụng các cấu hình ưu tiên an toàn như acks = all, enable.idempotence = true, enable.auto.commit = false, read_committed
  • Các vấn đề của Bufstream bao gồm việc consumer·producer bị dừng, phản hồi offset 0 sai, mất commit transaction, và mất dữ liệu ghi đã được xác nhận do lỗi lọc kích thước phản hồi của fetch API
  • Trong quá trình điều tra, cũng phát hiện ở Kafka Java client và Kafka transaction protocol các vấn đề như Consumer.close() bị block vô thời hạn, consumer offset không thể dự đoán, cùng các lỗi aborted read·lost write·torn transaction
  • Jepsen cho rằng do Kafka transaction protocol không bảo đảm tường minh thứ tự yêu cầu từ client và số hiệu transaction, nên tính an toàn của transaction trong Kafka và các hệ thống tương thích Kafka có thể bị phá vỡ khi dùng Java client chính thức

Cấu trúc Bufstream và phạm vi kiểm chứng

  • Kafka là một hệ thống streaming cung cấp log append-only có sao chép và sharding, còn Bufstream là một triển khai thay thế Kafka, ưu tiên quản trị dữ liệu và hiệu quả chi phí trong môi trường đám mây
  • Bufstream cung cấp topic và partition như Kafka, và hoạt động với Kafka client tiêu chuẩn
    • producer append record bằng producer.send()
    • consumer được gắn với partition bằng consumer.assign() hoặc consumer.subscribe(), sau đó đọc record bằng consumer.poll()
    • consumer group chia nhau xử lý record của một tập topic
  • Khi tích hợp với Buf Schema Registry, hệ thống có thể kiểm tra record Protocol Buffer để hỗ trợ xác thực record, kiểm soát truy cập ở mức field và chuyển đổi định dạng dữ liệu với các hệ thống khác
  • Khác với Kafka dùng đĩa cục bộ và giao thức sao chép riêng, Bufstream ghi dữ liệu trực tiếp vào object storage
    • nhằm tiết kiệm chi phí bằng cách tận dụng cấu trúc chi phí của lưu lượng sao chép trong object storage
    • node Bufstream có thể chạy như các VM stateless được auto-scale
  • Bufstream được cấu thành từ ba hệ con
    • agent: dịch vụ stateless cung cấp Kafka API
    • object store: lưu trữ các chunk record và cung cấp cho reader
    • coordination service: hiện dùng etcd, quyết định chunk nào đã được commit và thứ tự record
  • Tính đến tháng 10/2024, Bufstream mới chỉ được triển khai cho một số khách hàng; tài liệu quảng bá đây là “drop-in replacement” cho Apache Kafka và nhấn mạnh khả năng tương thích với Kafka transactions cùng exactly-once semantics, nhưng không có nhiều tuyên bố an toàn cụ thể

Cấu hình client và tiền đề transaction

  • Tương tự các bài kiểm thử trước đây với hệ thống tương thích Kafka, Jepsen đã điều chỉnh cấu hình client để có được hành vi an toàn hơn
  • Cấu hình Producer

    • sử dụng mặc định acks = all
    • trong Bufstream, acks = 0 có thể xác nhận ghi mà không chờ storage, dẫn tới khả năng mất dữ liệu ghi đã commit
    • acks = 1acks = all sẽ block cho đến khi Bufstream chắc chắn dữ liệu đã được lưu bền vững
    • để ngăn append trùng lặp do cơ chế retry tự động của Kafka producer, bài test dùng giá trị mặc định enable.idempotence = true
  • Cấu hình Consumer

    • do có tài liệu cảnh báo auto-commit có thể dẫn tới mất dữ liệu, nên nhìn chung sử dụng enable.auto.commit = false
    • khi không có committed offset, mặc định auto.offset.reset sẽ bắt đầu từ offset mới nhất, nên không bảo đảm at-least-once delivery
    • để consumer có thể quan sát toàn bộ log, sử dụng auto.offset.reset = earliest
    • Kafka transaction được cấu thành từ tập record do producer gửi và bản đồ offset tối đa theo từng partition mà consumer đã poll
    • chỉ khi transaction được commit thì các record đã gửi mới bền vững, cuối cùng mới hiển thị với consumer read_committed, và committed offset cũng tăng lên đến ít nhất offset được chỉ định trong transaction
    • nếu transaction không được commit thì committed offset không tiến lên, còn khả năng nhìn thấy dữ liệu ghi có thể thay đổi tùy cấu hình consumer
    • hiện tượng consumer read_uncommitted đọc được giá trị của transaction đã abort được phân loại là aborted read(G1a)
    • tài liệu Kafka nói rằng read_committed ngăn G1a và phần nào bảo đảm tính chất mọi ghi của transaction либо đều hiển thị hoặc không hiển thị gì cả, nhưng trong các bài kiểm thử Kafka·Redpanda·Bufstream của Jepsen vẫn quan sát được write cycle (hiện tượng tương tự G0) và một số dạng G1c

Thiết kế kiểm thử

  • Jepsen đã kiểm thử Bufstream từ 0.1.0 đến 0.1.3 cùng nhiều bản release candidate build
  • Test harness sử dụng Bufstream test harness, Jepsen testing library, và Java Kafka Client 3.8.0
  • Môi trường chạy

    • sử dụng 3~5 node Debian Bookworm trên cả LXC container và EC2 VM
    • 1 node cho etcd, 1 node cho Minio, các node còn lại dùng làm Bufstream agent
    • producer, consumer và admin client được khởi tạo với chỉ một node duy nhất trong bootstrap_servers, nhưng không chặn smart client discovery
  • Các thiết lập an toàn chính

    • auto-commit false
    • acks = all
    • retries 1,000
    • idempotence enabled
    • isolation level read_committed
    • auto_offset_reset = earliest
    • tắt tự động tạo topic phía server
    • tiêm lỗi bao gồm process pause(SIGSTOP), crash(SIGKILL), clock skew(clock_settime), network partition(iptables)
    • do Bufstream tách thành agent, object store và coordination service, Jepsen đã tạo mới công cụ để tiêm lỗi nhắm vào từng hệ con cụ thể
    • ví dụ có thể crash chỉ các node Bufstream hoặc chỉ pause coordinator etcd, rồi thay đổi các tổ hợp này theo thời gian

Queue workload và Abort workload

  • Queue workload phân tích tính an toàn theo mô hình dữ liệu Kafka
    • Mỗi logical process chạy producer, consumer và admin client
    • Numeric key dùng để nhận diện một topic-partition cụ thể
    • Key được chọn theo tần suất hàm mũ, nên một số key được truy cập thường xuyên còn một số khác thì hiếm hơn
  • Sử dụng ba operation cơ bản
    • crash: kết thúc logical process và thay thế bằng client mới
    • subscribe hoặc assign: thay đổi tập topic hoặc partition mà consumer sẽ poll
    • txn, poll, send: thực hiện một chuỗi micro-operation poll hoặc send
  • Trong workload non-transactional, mỗi send hoặc poll chỉ chứa đúng một micro-operation
  • Trong workload transactional, nhiều micro-operation được bọc trong một Kafka transaction
  • Phân tích tạo mapping offset-to-value theo từng key rồi tìm lỗi
    • Nếu có nhiều value xuất hiện tại cùng một offset thì là inconsistent offset
    • Nếu cùng một value xuất hiện ở nhiều offset thì là duplicate error
    • Nếu record đã được chấp nhận nhưng hoàn toàn không được quan sát thấy thì là lost hoặc unseen
    • Nếu poll trả về value được gửi bởi operation đã abort thì là aborted read
    • Cũng kiểm tra liệu transaction có quan sát được chính bản ghi của nó hay không
  • Sau bài kiểm thử chính, hệ thống sẽ được khôi phục lỗi và chuyển sang giai đoạn final reads
    • Mỗi process đọc mọi topic-partition từ offset 0 và poll đến offset đã ghi cao nhất được biết đến
    • Nếu final reads bị timeout và record đã được chấp nhận vẫn không được quan sát thấy, nó sẽ được phân loại là unseen
  • Abort workload được bổ sung để theo dõi hành vi poll offset sau khi transaction abort
    • Topic bị giới hạn ở một partition, process, producer và consumer duy nhất
    • Sau khi transaction poll record rồi cố ý abort, poll offset tiếp theo được phân loại thành advance, rewind, rewind-further hoặc other

5 vấn đề được phát hiện trong Bufstream

  • Consumer bị kẹt (#1)

    • Từ 0.1.0 đến 0.1.3-rc.8, giai đoạn đọc cuối cùng thường xuyên bị kẹt
    • consumer.poll() trả về kết quả rỗng ngay lập tức, nhưng trong log vẫn còn lại hàng nghìn record đã được xác nhận
    • Trạng thái này kéo dài từ vài chục giây đến hơn 1 giờ
    • Trong một bài test, 691 record đã được xác nhận được gửi trong 120 giây đầu tiên, và tại thời điểm bắt đầu final reads, có 40 record không được bất kỳ poller nào quan sát thấy
    • Sau đó, consumer.poll() không trả về kết quả trong hơn 1 giờ nên bài test bị timeout
    • Nguyên nhân là Bufstream node sau khi khởi động lại có thể trả về giá trị cache cũ của last stable offset và high watermark
    • Một số client library kết luận rằng không còn record nào ở phía sau và bị stall; Bufstream đã áp dụng bản vá trong 0.1.3-rc.6 để refresh cache khi khởi động
  • Producer và consumer bị kẹt (#2)

    • Ngay cả trong 0.1.3-rc.6, vấn đề unseen write vẫn tiếp tục được quan sát sau pause, crash và partition đối với coordinator, storage và Bufstream node
    • Trong một số trường hợp, sau khi coordinator pause, dù tất cả Bufstream node đều đang chạy, client vẫn rơi vào trạng thái chờ InitProducerId rồi timeout
    • Trong các trường hợp khác, listOffsets thất bại với node ... being disconnected hoặc timed out waiting for a node assignment, còn poll thì hoàn tất nhưng không trả về kết quả
    • Kill rồi restart Bufstream node sẽ giải quyết được vấn đề
    • Nguyên nhân liên quan đến etcd lease
    • Bufstream agent dùng etcd leases để theo dõi active agent
    • Do pause hoặc partition ngắn, etcd đã xóa key gắn với agent lease, nhưng update xóa có thể không được gửi tới agent
    • Agent rơi vào trạng thái không biết rằng nó đã mất lease của chính mình
    • Nhóm Bufstream đã thêm logic polling bổ sung, và trong 0.1.3-rc.8, unseen write nhìn chung đã được khắc phục
  • Offset 0 giả (#3)

    • Từ 0.1.0 đến 0.1.3-rc.2, giá trị đã gửi có thể được gán offset 0 rồi xuất hiện ở offset thực tế cao hơn
    • Điều này vẫn xảy ra ngay cả khi offset 0 đã được gán từ rất lâu trước đó
    • Chỉ sender quan sát thấy offset 0, còn poller thì quan sát thấy offset cao hơn
    • Trong một bài test 2 phút với một Bufstream node và pause tiến trình etcd, 6 write nhận offset 0 rồi sau đó xuất hiện ở offset cao hơn
    • Nguyên nhân là phản hồi lỗi của Bufstream thiếu field bắt buộc
    • Bufstream gửi yêu cầu commit log tới etcd và etcd đã xử lý, nhưng do pause hoặc partition, Bufstream có thể timeout trong lúc chờ response
    • Bufstream gửi error code cho client, nhưng không đặt offset của record đã gửi thành -1, tức tín hiệu lỗi
    • Java Kafka client diễn giải điều này là phản hồi thành công với offset 0
    • Franz-go, được Bufstream test suite sử dụng, diễn giải message này là lỗi nên vấn đề không lộ ra trong bài test
    • Bufstream đã sửa trong 0.1.3-rc.6, và sau đó Jepsen không thể quan sát lại nữa
  • Mất write của transaction (#4)

    • Trong 0.1.2, tình trạng mất write khiến một số record của transaction đã commit biến mất và không bao giờ được quan sát lại xảy ra thường xuyên
    • Trong một bài test, trong 100 giây và 6.761 write transaction, 240 record do các transaction đã commit ghi ra đã bị mất
    • Trong ví dụ, value 141 của key 5 được trả về là đã ghi thành công ở offset 274, nhưng mọi consumer.poll() đều bỏ qua offset đó
    • Nguyên nhân là bug trong cơ chế an toàn đồng thời được thêm vào 0.1.2
    • Cơ chế này gán số duy nhất cho mỗi transaction trong producer epoch để giảm bớt sự thiếu idempotence của Kafka transaction protocol
    • Do bug trong logic theo dõi transaction number, khi nhiều transaction được commit qua nhiều epoch, một số commit bị bỏ qua nhầm
    • Transaction trông như đã được commit thực tế có thể bị abort, hoặc ngược lại
    • Jepsen phát hiện bug này nhờ đặt transaction timeout thấp ở mức 1 giây
    • Bufstream đã xác định vấn đề chỉ vài giờ sau khi phát hành 0.1.2, chặn khách hàng nâng cấp, và khách hàng đã không nâng cấp lên 0.1.2
    • Bản sửa nằm trong 0.1.3-rc2
  • Lost writes do server-side filtering (#5)

    • Trong 0.1.3-rc.8, sau các lỗi nhỏ như pause tiến trình Bufstream hoặc coordinator, hay partition giữa hai bên, thường xuyên xuất hiện một khoảng thời gian ngắn xảy ra write loss
    • Mất dữ liệu xảy ra bất kể có dùng transaction hay không
    • Trong một bài test 5 phút, trong số 16.770 record có 22 record đã được acknowledge nhưng không consumer nào poll được
    • Một số record có lúc được poller nhìn thấy, nhưng về sau lại biến mất khỏi poll
    • Nguyên nhân là logic giới hạn kích thước response của fetch API được thêm vào 0.1.3-rc.8 để workaround bug của một Kafka web GUI phổ biến
    • Bug trong logic filtering đã ẩn record khỏi consumer đang bị lag, khiến nó trông như write loss
    • Bufstream đã sửa trong 0.1.3-rc.12

Vấn đề với Kafka Java client và giao thức Kafka

  • KIP-588: ProducerFencedException dễ gây hiểu nhầm

    • Trong quá trình kiểm thử, lỗi ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one. xuất hiện thường xuyên
    • Ngay cả trong bài kiểm thử mà mọi producer đều nhận transactional ID duy nhất, lỗi này vẫn xuất hiện nên mất khá nhiều thời gian để xác định nguyên nhân
    • KIP-588 ghi rằng ngay cả khi transaction timeout thì ProducerFencedException cũng có thể được ném ra
    • Kafka Java client dùng TimeoutException chuyên biệt cho hầu hết các trường hợp timeout, nhưng trong trường hợp này lại ném ProducerFencedException
    • Dù trên thực tế không có producer xung đột nào, thông báo lỗi vẫn nói rằng tồn tại một instance producer thứ hai
    • KIP-588 đã được mở suốt 2 năm, và Jepsen khuyến nghị nhóm Kafka thay đổi thông báo lỗi
  • KAFKA-17734: Consumer.close() có thể block vô thời hạn

    • Ở cả bài kiểm thử với Bufstream lẫn Kafka, quá trình kiểm thử bị dừng lại sau vài giờ do bug của Java client
    • Consumer.close() mặc định sẽ block ở network IO
    • Tham số timeout của close() lẽ ra phải ngăn việc block vô thời hạn, nhưng lại không hoạt động
    • Cách gọi consumer.wakeup() từ một thread riêng để interrupt consumer đang bị kẹt ở IO cũng không hiệu quả
    • Jepsen cho rằng các chương trình chạy dài hạn phải có khả năng giải phóng các resource như client, connection, thread, memory trong thời gian hợp lý ngay cả khi có lỗi mạng, và đã đăng KAFKA-17734
  • KAFKA-17582: sau khi transaction thất bại, consumer offset trở nên khó dự đoán

    • Tài liệu chính thức của Kafka gần như không nói rõ consumer offset nên ở trạng thái nào khi transaction commit thất bại
    • Tài liệu thiết kế Kafka của Confluent nói rằng nếu transaction bị abort thì consumer position sẽ quay lại giá trị trước đó, nhưng Java client thực tế không phải lúc nào cũng hoạt động như vậy
    • Kết quả của workload abort cho thấy ngay cả trên một cluster khỏe mạnh, hành vi sau abort cũng rất khó dự đoán
    • Phần lớn các cặp transaction đều advance tới offset xa hơn
    • Một số lại rewind về offset trước đó
    • Mọi lần rewind đều liên quan đến sự kiện rebalance, và mọi lần advance đều không có rebalance
    • Theo phản hồi từ phía Kafka, hành vi này là có chủ đích
    • consumer sẽ tiếp tục advance
    • Khi rebalance xảy ra, nó có thể bị rewind về một điểm bất kỳ dựa trên committed offset
    • Người dùng phải tự tay rewind consumer position khi transaction bị abort
    • Jepsen đã mở KAFKA-17582, đề xuất tài liệu hóa hành vi này và cân nhắc thay đổi mặc định sang rewind khi transaction abort
    • Queue workload cũng đã được sửa để rewind consumer một cách tường minh
  • KAFKA-17754: mất ghi, đọc dữ liệu đã abort, torn transaction

    • Trên Bufstream 0.1.0~0.1.3, chỉ với việc pause tiến trình Bufstream, pause coordinator, crash, hoặc network partition, đã quan sát được aborted read, lost write và vi phạm tính nguyên tử
    • Phân tích dẫn đến một khiếm khuyết mang tính nền tảng trong transaction protocol của Kafka
    • Trong ví dụ, client thực thi transaction với transactional ID duy nhất jt1234 và gửi committed = false trong EndTxn để abort, nhưng 15 lần gọi poll() vẫn quan sát được các bản ghi do transaction đã abort ghi ra
    • Những bản ghi khác của cùng transaction thì không poller nào quan sát được
    • Khi đối chiếu packet capture với log của Bufstream, nguyên nhân được xác định là một commit message bị xử lý trễ
    • Một EndTxn commit được gửi từ vài transaction trước đã bị một node xử lý muộn
    • Trong khi đó, client đã tiếp tục thực hiện các transaction tiếp theo
    • Commit bị trễ đó đã được áp dụng lên transaction hiện tại, khiến chỉ phần đầu của transaction được commit, còn phần còn lại bị xử lý như transaction riêng và bị abort
    • Kafka protocol được thiết kế để client có thể gửi request qua nhiều TCP connection và tới nhiều node, nhưng không có sequence number để xác định thứ tự request của cùng một client
    • Cũng không có khái niệm transaction number, nên khi server nhận commit hoặc abort message, nó không thể biết client đang muốn kết thúc transaction nào
    • Kết quả là có thể xảy ra các tình huống sau
      • transaction trông như đã commit nhưng thực tế lại bị abort
      • transaction đã abort nhưng thực tế lại commit
      • chỉ một phần write của transaction được giữ lại, phần còn lại bị mất, tạo ra torn transaction
    • Kafka Java client chính thức coi timeout là retryable và có thể tự động gửi nhiều EndTxn message, nên ngay cả khi người dùng chỉ gọi commit hoặc abort đúng một lần cho mỗi transaction thì vấn đề vẫn có thể phát sinh
    • Jepsen cũng quan sát thấy aborted read và torn transaction trong Kafka khi pause tiến trình, và đã mở KAFKA-17754
    • Các kỹ sư Kafka cho rằng KIP-890 có thể khắc phục vấn đề này
    • KIP-890 thay đổi transaction protocol theo hướng tăng producer epoch cho mỗi transaction
    • Vì server từ chối các message có epoch cũ, nó có thể ngăn commit message của transaction trước bị rò sang transaction sau
    • Trong bản 0.1.3, Bufstream đã thêm một mechanism dùng etcd revision làm logical clock để giảm tần suất xảy ra, nhưng vẫn không thể ngăn việc reorder giữa client và Bufstream
    • Jepsen tiếp tục quan sát thấy aborted read, lost write và torn transaction cả trên 0.1.3, và cho rằng cần có giải pháp ở phía client

Tóm tắt kết quả tổng thể

  • Cả 5 vấn đề nội tại của Bufstream đều đã được khắc phục
    • #1: consumer bị kẹt do highest stable offset bị trễ, không cần gián đoạn, đã sửa trong 0.1.3-rc.6
    • #2: producer/consumer bị kẹt do etcd lease hết hạn, cần pause, đã sửa trong 0.1.3-rc.8
    • #3: xuất hiện zero offset giả, cần pause, đã sửa trong 0.1.3-rc.6
    • #4: mất các bản ghi ghi của transaction, không cần gián đoạn, đã sửa trong 0.1.3-rc.2
    • #5: mất dữ liệu ghi do server-side filtering, cần pause, đã sửa trong 0.1.3-rc.12
  • Các vấn đề liên quan đến Kafka vẫn còn tồn tại
    • KIP-588: thông báo lỗi sai ở transaction timeout, chưa được giải quyết
    • KAFKA-17734: ConsumerClient.close() có thể block vô thời hạn, chưa được giải quyết
    • KAFKA-17582: sau khi transaction thất bại, consumer offset trở nên khó dự đoán, chưa được giải quyết
    • KAFKA-17754: mất dữ liệu ghi, aborted read, torn transaction, chưa được giải quyết
  • Jepsen lưu ý rằng việc kiểm chứng an toàn mang tính thực nghiệm có thể chứng minh bug tồn tại, nhưng không thể chứng minh bug không tồn tại
  • Đặc biệt, Jepsen cho rằng KAFKA-17754 khiến việc xác định Bufstream còn có các trường hợp mất dữ liệu ghi khác hay không trở nên khó khăn

Khuyến nghị cho người dùng và vận hành Bufstream

  • Người dùng sử dụng transaction của Bufstream với Java Kafka client chính thức hiện cần cân nhắc rằng transaction có thể không an toàn
    • transaction đã abort có thể trên thực tế lại được commit
    • transaction đã commit có thể trên thực tế lại bị abort
    • transaction có thể bị xé đôi, chỉ giữ lại một phần hiệu ứng
  • Bufstream cho rằng Franz-go client ít dễ dính vấn đề này hơn, nhưng Jepsen chưa kiểm thử Franz-go bằng các kỹ thuật như trong công việc này
  • Các client khác có thể bị ảnh hưởng hoặc không
  • Người dùng Bufstream trước 0.1.3 có thể gặp các vấn đề sau
    • producer.send() trả về sai offset 0 thay vì offset thực tế
    • vấn đề khả dụng metastable khiến client bị kẹt
  • Jepsen khuyến nghị nâng cấp lên 0.1.3
  • Jepsen đánh giá kiến trúc tổng thể của Bufstream có vẻ sound
    • cách dùng một coordination service như etcd để xác định thứ tự của các immutable data chunk là một cách tiếp cận tương đối đơn giản, đã có tiền lệ trong OLTP và các hệ thống streaming
  • Về mặt vận hành, có hai cải tiến được khuyến nghị
    • nếu yêu cầu shared file của storage thất bại lúc startup thì cluster có thể crash, nên Jepsen khuyến nghị thêm retry và Bufstream đã bổ sung một lớp retry
    • khi dependency không khả dụng, nên để agent tiếp tục chạy thay vì chết ngay, đồng thời cung cấp backpressure và trạng thái hệ thống để phục hồi mượt hơn
  • Tính đến 0.1.3, Bufstream đã thêm logic retry bổ sung cho etcd, nhưng vẫn cần constant supervision để duy trì trạng thái online
  • Người dùng nên bảo đảm có process supervisor và kiểm thử xem hệ thống có tiếp tục hoạt động mà không bỏ cuộc ngay cả trong các đợt outage kéo dài hay không

Cần tài liệu hóa và sửa đổi giao thức transaction của Kafka

  • Tài liệu chính thức của Kafka hầu như không nói gì về transaction, khiến người dùng phải ghép nhiều nguồn mơ hồ và đôi khi mâu thuẫn với nhau
  • Jepsen đã khuyến nghị đội Kafka tạo một tài liệu trung tâm làm rõ transaction semantics và nhắc đến KAFKA-17671
  • Tài liệu đó tối thiểu cần nêu rõ những điểm sau
    • khi nào consumer quan sát được offset tăng đơn điệu
    • khi nào consumer có thể bỏ qua các record đã được acknowledge
    • rebalance có thể ảnh hưởng giữa chừng của transaction hay không
    • khi nào write offset của producer tăng đơn điệu
    • khi nào G0, G1a, G1b, G1c, fractured read và việc đọc các bản ghi do chính transaction của mình ghi ra là hợp lệ
    • sau một transaction bị abort, giá trị trả về của poll() và offset mang ý nghĩa gì
    • cần xử lý ra sao với transaction error, lỗi trong lúc abort, và lỗi trong lúc rewind
  • Jepsen chỉ ra rằng tài liệu của Confluent lặp đi lặp lại việc mặc định của Kafka cung cấp at-least-once delivery, nhưng điều này có vẻ không đúng
    • auto.offset.reset = latest có thể khiến các record chưa được xử lý trông như đã được “committed”
    • tài liệu về offset management của Confluent cũng nói đến rủi ro mất tiến độ xử lý message khi crash dưới cơ chế auto-commit mặc định
    • tài liệu nói consumer sẽ bị rewind khi transaction abort cũng không khớp với thực tế
  • Jepsen cho rằng giao thức transaction của Kafka cần được sửa đổi ở cấp độ nền tảng
    • giao thức ngầm giả định ordered reliable delivery, nhưng trên thực tế tồn tại process pause, mạng không đáng tin cậy, độ trễ khác 0 và việc phân phối không theo thứ tự giữa nhiều TCP socket
    • giao thức Kafka phân tán message qua nhiều node và TCP socket, còn client thì tự động retry message
    • không có sequence number để khôi phục thứ tự của cùng một client message, cũng không có transaction number để xác nhận transaction đích
  • KIP-890 cố gắng bảo đảm thứ tự chặt chẽ hơn bằng cách tăng epoch ở mỗi lần commit transaction
  • Thư viện client cũng có thể hỗ trợ bằng cách re-initialize producer để tăng epoch khi message không được acknowledge
  • Java Kafka Client 3.8.0 dễ bị ảnh hưởng bởi vấn đề này
  • Jepsen cho rằng Franz-go có thể giảm nhẹ hoặc ngăn vấn đề nhờ thực hiện re-initialize khi timeout, nhưng họ chưa điều tra các thư viện client khác

Công việc sắp tới

  • Nhiều người dùng dựa vào “exactly-once semantics” của Kafka Streams API hơn là trực tiếp xử lý transaction, nên trong tương lai có thể khảo sát tính chính xác của ứng dụng Streams
  • Jepsen cũng gặp unseen write trong Kafka khi điều tra KAFKA-17754, nhưng do hạn chế về thời gian nên chưa thể phân tích
    • unseen write có thể là dấu hiệu của hanging transaction, stuck consumer hoặc mất dữ liệu
    • Vẫn còn câu hỏi liệu message Produce bị trì hoãn có thể đi vào transaction trong tương lai và vi phạm transaction guarantee hay không
    • Cũng có nghi ngờ rằng Kafka Java Client tái sử dụng sequence number khi request timeout, khiến write đã được acknowledge nhưng có thể bị loại bỏ âm thầm
  • Khi xảy ra rebalance event, consumer position có thể di chuyển tới lui, nhưng quy tắc của việc này vẫn chưa rõ ràng
  • Nếu Kafka tài liệu hóa hành vi dự kiến, Jepsen muốn xác minh điều đó
  • Jepsen giải thích rằng đây là một quy trình ngẫu nhiên nên khó khám phá các anomaly hiếm gặp
    • Những vấn đề chỉ xảy ra một lần thì việc debugging và tái hiện là cực kỳ khó
  • Bufstream cũng sử dụng Antithesis, nơi chạy toàn bộ hệ thống phân tán trong hypervisor xác định và mạng mô phỏng
    • Kết hợp workload generation và history checking của Jepsen với môi trường deterministic, có thể replay của Antithesis có thể cải thiện khả năng tái hiện bài kiểm thử

1 bình luận

 
GN⁺ 2024-11-14
Các ý kiến trên Hacker News
  • Nếu trong lúc điều tra các vấn đề như KAFKA-17754 mà còn phát hiện cả ghi vô hình trong Kafka, có lẽ đã đến lúc Jepsen đào sâu lại Kafka
    Lần điều tra cuối là năm 2013 (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta), còn hiện tại trông như mới ở giai đoạn bắt đầu phát hiện nhiều vấn đề trong chính Kafka
    Những chuyện kiểu “ghi đã được xác nhận nhưng có thể bị âm thầm vứt bỏ” nghe khá đáng sợ

    • Tôi rất muốn phân tích Kafka :-)
  • Phần nói rằng với giá trị mặc định enable.auto.commit=true, consumer Kafka có thể commit offset bất kể ứng dụng đã thật sự xử lý hay chưa là rất đáng ngạc nhiên
    Tôi chưa từng hiểu auto commit theo cách đó, và nếu đó là mặc định thì thấy thật vô lý
    Phần mô tả trong tài liệu không thật rõ ràng, nhưng nhìn tổng thể tôi đọc ra rằng offset chỉ được commit khi xử lý đã hoàn tất
    Tôi hiểu việc điều chỉnh khoảng thời gian auto commit, như kỳ vọng trong xử lý ít nhất một lần (at-least-once), là giúp giảm cửa sổ xử lý trùng lặp chứ không phải gây mất message

    • Hơi bất ngờ, và tôi đồng ý rằng tài liệu giải thích phần này chưa tốt
      Nếu không commit một cách tường minh, Kafka không có cách nào biết message đã được xử lý hay chưa
      Kafka giả định rằng message đã giao được xử lý ngay lập tức
      Auto commit giống như đưa một cây kem ốc quế rồi quay đi ngay và giả định người kia đã ăn nó. Có người có thể vừa nhận đã làm rơi và chưa kịp ăn miếng nào
    • Điểm mấu chốt là message được chuyển thành công tới client Kafka không có nghĩa là ứng dụng đã xử lý nó
      Nếu muốn bảo đảm điều đó, bạn phải acknowledge một cách tường minh
      Ví dụ nếu việc cần làm chỉ là ghi message vào database, thì ngay khi message đi vào callback handler của client, nó được coi là đã xác nhận
      Nhưng trên thực tế, rất có thể bạn muốn nó chỉ được xác nhận sau khi insert vào DB thành công
      Nếu DB không truy cập được vì mạng, Kubernetes, cấu hình firewall, v.v., và trong lúc đó kỹ sư thử khởi động lại làm client bị tắt, rất dễ có message chưa được xử lý
    • Tôi hiểu tính năng này dành cho các tình huống hiệu năng cao
      Một hệ thống khác có thể xác định việc thất bại, và tính năng này có thể dịch chuyển vị trí giới hạn trên để giảm việc xử lý lại
      Tuy nhiên nếu đúng thời điểm và có lỗi xảy ra, phải giả định rằng sau khi khởi động lại bạn có thể nhận lại một phần đã xử lý
      Vấn đề là khi trước auto commit không có kiểu xử lý như vậy
      Đọc thì có vẻ nó được thiết kế để commit khá lâu sau khi xử lý, nhưng việc vừa là auto commit vừa chỉ nên commit các mục cách thời điểm auto commit vài mili giây nghe cũng có phần mâu thuẫn
    • Có thể biện minh phần nào cho lý do tính năng này tồn tại. Nó được thiết kế cho consumer đồng bộ, đơn luồng, và giả định một vòng lặp đại khái là gọi poll rồi xử lý message một cách bền vững
      Điểm dễ gây nhầm lẫn là việc kiểm tra auto commit không diễn ra bất đồng bộ sau timeout, mà diễn ra ở lần gọi poll tiếp theo
      Vì vậy chỉ khi bạn không xử lý message một cách bền vững trước khi gọi lại poll mà chỉ lưu nó lại, chẳng hạn dùng xử lý bất đồng bộ, trì hoãn, queue, v.v., thì mới có khả năng làm rơi ghi
      Điều này dựa trên hành vi đã được tài liệu hóa của thư viện client Java (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...); việc implementation hiện tại có thật sự như vậy hay không là chuyện khác
      Protocol Kafka bị kẹt giữa mức cao và mức thấp nên cả hai phía đều không thật sự tốt
      Auto commit là tính năng mức cao giúp tạo ứng dụng đơn giản dễ hơn, nhưng nếu không dùng theo cách nó kỳ vọng thì dĩ nhiên có thể thất bại
      Tôi nghĩ ngày nay người dùng cuối nên dùng các implementation mức cao xử lý đúng các chi tiết này hơn là trực tiếp dùng client Kafka. Nếu dùng cho dữ liệu thì là stream processing engine, còn cho ứng dụng thì là một persistent execution engine
  • Nhìn vào trang sản phẩm (https://buf.build/product/bufstream), tôi thắc mắc làm sao mô tả “chỉ chạy trong VPC AWS hoặc GCP của bạn và không liên lạc ra bên ngoài” có thể đi cùng tính phí theo mức sử dụng “$0.002 mỗi GiB trước nén”
    Chắc họ không vận hành cả mô hình kinh doanh bằng hệ thống danh dự đâu

    • Phần giới thiệu nói “tính đến tháng 10/2024, Bufstream chỉ được triển khai cho các khách hàng được chọn”, nên tôi nghĩ hệ thống danh dự cũng có thể khả thi
      Tất nhiên có rủi ro bị lạm dụng, nhưng đó có thể là một đánh đổi đáng giá để thu hút một số khách hàng nhất định
    • Chương trình hoặc là mã nguồn mở, hoặc là không
      Nếu source không được công khai, tuyệt đối không nên tin tuyên bố “không liên lạc ra bên ngoài”
  • “Protocol transaction của Kafka về căn bản đã hỏng và cần được sửa đổi” nghe đau thật
    Dù vậy, như mọi khi, phần điều tra và bài viết rất xuất sắc

  • Không biết Kyle đã từng xem xét NATS JetStream chưa. Tôi tò mò anh ấy sẽ nghĩ gì

    • Chưa xem xét, nhưng bạn không phải người đầu tiên yêu cầu
      Một số người đã gợi ý rằng nó sẽ… nói sao nhỉ… thú vị :-)
  • Tôi không tìm thấy dự án bufstream trên GitHub, không biết nó ở đâu

  • Sau khi đọc bài blog và tài liệu liên quan, có vẻ “exactly-once delivery” của Kafka được định nghĩa là thuộc tính của tác vụ đọc-xử lý-ghi, trong đó worker đọc từ topic 1 và ghi vào topic 2, với cả hai topic nằm trong cùng một hệ thống Kafka logic
    Nếu đúng vậy thì tôi nghĩ gọi nó là transaction sẽ hợp hơn

    • Kafka thực ra cũng gọi nó là transaction
      Tuy nhiên có hai cách nhìn về “exactly once”
      Một là theo nghĩa giống transaction database: hiệu ứng không được bị lặp hoặc biến mất
      Cách còn lại giống một thuộc tính của đồ thị luồng dữ liệu về quan hệ giữa các message băng qua các topic-partition, gần với consistency trong ACID hơn một chút
      Tương tự như một hệ thống transaction serializable bảo đảm một tính nhất quán ở cấp miền nhất định, bạn có thể dùng transaction để đạt được thuộc tính luồng dữ liệu đó
      Ví dụ, serializability bảo đảm rằng các invariant được bảo toàn khi xét từng transaction riêng lẻ cũng được bảo toàn trong lịch sử thực thi đồng thời
      Có thể xem Kafka đang cố đạt tới “exactly-once semantics” theo cách đó
  • Đừng nhầm với https://www.warpstream.com/

    • Đúng vậy. WarpStream thậm chí không hỗ trợ transaction
  • Errata: “Transactions may observe none, part, or all” có lẽ nên là “Consumers may observe none, part, or all”

    • Cả hai đều đúng, nhưng tôi viết là transaction để rõ ràng hơn
      Semantics của consumer bên ngoài transaction còn mơ hồ hơn
      Tất cả các lần đọc trong workload này đều diễn ra trong ngữ cảnh transaction và đi qua đường commit offset transaction
  • Tôi tò mò phần mềm này được dùng vào đâu. Instrumentation? Hộp đen?

    • Jepsen là công cụ khiến bạn muốn khóc nếu không biết rằng nó đang kiểm thử database do bạn phát triển
      Dĩ nhiên là nước mắt vui mừng. Vì việc được Jepsen chú ý tự nó đã là một thành tựu
    • Đây là clone của Kafka. Kafka về cơ bản là một queue bền vững