Xây dựng lại hệ thống phát hiện sự cố với Kafka, Flink và OpenTelemetry
Một nhóm kỹ sư đã giảm độ trễ phát hiện sự cố từ 40 giây xuống dưới 10 giây bằng cách thay thế trình tổng hợp Node.js bằng Apache Flink trên Kubernetes và OpenTelemetry.
Được dịch tự động từ bản gốc tiếng Anh.
Một nhóm kỹ sư nhỏ đã xây dựng lại nền tảng phát hiện sự cố tự động của họ, cắt giảm độ trễ từ sự kiện đến chỉ số (event-to-metric) từ hơn 40 giây xuống còn dưới 10 giây. Kiến trúc mới dựa vào Apache Kafka, Apache Flink chạy trên Kubernetes và OpenTelemetry để xử lý hàng tỷ sự kiện vận hành phía khách hàng mỗi ngày. Mặc dù hệ thống cải thiện tốc độ và tính cách ly, nhóm báo cáo rằng tỷ lệ thu hồi (recall rate) dao động từ 64% đến 86%, nhấn mạnh sự phức tạp liên tục trong việc phát hiện bất thường chính xác.
Chuyện gì đã xảy ra
Tổ chức này vận hành hơn mười sản phẩm đám mây phục vụ hàng triệu khách thuê bao (tenant). Trước đây, ngăn xếp giám sát của họ sử dụng một trình tổng hợp Node.js được cung cấp dữ liệu bởi một hàng đợi đám mây dùng chung, chạy trên khoảng 90 máy ảo. Hệ thống cũ này gặp phải độ trễ cao, vấn đề "hàng xóm ồn ào" (noisy neighbor) khi lưu lượng truy cập của các khách thuê bao khác gây ra độ trễ, và chi phí tăng tuyến tính với mỗi sản phẩm mới được tích hợp. Chi phí vận hành hàng năm đã tăng từ 120.000 USD lên 230.000 USD, và tầng bộ nhớ đệm thường xuyên đạt mức CPU 100% trong các thay đổi thông thường.
Để giải quyết những hạn chế này, nhóm đã thiết kế một đường ống mới tập trung vào năm mục tiêu: độ trễ dưới 10 giây, đường dẫn tiêu thụ riêng biệt cho tính cách ly, ghi dữ liệu idempotent (không phụ thuộc vào trạng thái trước đó) để đảm bảo tính đúng đắn khi phát lại, chi phí mở rộng theo khối lượng thay vì số lượng tính năng, và khả năng vận hành dựa trên cấu hình. Kết quả là một hệ thống mà việc tích hợp một trải nghiệm mới chỉ yêu cầu một pull request để cập nhật cấu hình, không cần triển khai toàn bộ. Nhóm đã đo lường hiệu suất theo từng tháng trong mười tám tháng, ghi nhận rằng mặc dù tốc độ cải thiện đáng kể, độ chính xác vẫn thấp hơn mục tiêu của họ.
Cách thức hoạt động
Lớp upstream sử dụng Apache Kafka làm bus sự kiện. Thay vì tiêu thụ toàn bộ luồng dữ liệu khổng lồ (firehose), nhóm đã triển khai một bộ lọc đăng ký phía máy chủ được định nghĩa trong mã nguồn. Bộ lọc này chỉ cho phép các sản phẩm và trải nghiệm cụ thể, đồng thời loại bỏ lưu lượng thử nghiệm và tổng hợp. Dữ liệu đã lọc được đưa vào một topic Kafka chuyên dụng với thời gian lưu giữ bảy ngày, đóng vai trò là cửa sổ phát lại để gỡ lỗi hoặc khôi phục. Hai đầu vào bên cạnh (side inputs) cung cấp dữ liệu cho đường ống: một dịch vụ ngữ cảnh khách thuê bao cung cấp siêu dữ liệu như shard và khu vực, và một kho lưu trữ cấu hình.
Ở phần lõi là một job Apache Flink 1.20 duy nhất được triển khai thông qua Flink Kubernetes Operator. Job này lọc bỏ các sự kiện cũ và lỗi, làm phong phú dữ liệu bằng ngữ cảnh khách thuê bao sử dụng một sidecar bất đồng bộ với các circuit breaker, và tổng hợp các chỉ số. Nó sử dụng các phác thảo HyperLogLog để ước tính số người dùng bị ảnh hưởng riêng biệt mà không đếm trùng lặp, đạt biên độ lỗi khoảng 1,5%. Trạng thái được quản lý trong RocksDB với các checkpoint 30 giây tới lưu trữ đối tượng, đảm bảo ngữ nghĩa exactly-once cho các tệp Parquet và ghi idempotent vào một kho lưu trữ key-value. Các chỉ số được xuất qua OpenTelemetry tới một cơ sở dữ liệu chuỗi thời gian tương thích Prometheus.
Mặt phẳng quyết định, có tên là AutoHOT, chạy dưới dạng một dịch vụ Go ở hai khu vực. Nó tiêu thụ cảnh báo từ các trình phát hiện, định lượng tác động bằng cách truy vấn kho dữ liệu tổng hợp và áp dụng ma trận mức độ nghiêm trọng. Để ngăn chặn dương tính giả, nó kiểm tra các biến động đột ngột trong hoạt động của người dùng trong mười lăm phút trước khi tạo ticket. Engine sử dụng khóa phân tán để đảm bảo chỉ một khu vực xử lý một cảnh báo tại một thời điểm, ngăn ngừa sự cố trùng lặp trong quá trình chuyển đổi dự phòng (failover). Nó cũng giám sát sự im lặng trong telemetry, nhận ra rằng thiếu dữ liệu có thể chỉ ra một shard cơ sở dữ liệu bị sập hoàn toàn.
Chi tiết quan trọng
- Độ trễ giảm từ hơn 40 giây xuống dưới 10 giây cho quy trình xử lý từ sự kiện đến chỉ số.
- Hệ thống xử lý hàng tỷ sự kiện mỗi ngày bằng một job Apache Flink 1.20 duy nhất trên Kubernetes.
- Các phác thảo HyperLogLog cho phép đếm số người dùng riêng biệt có thể hợp nhất giữa các khách thuê bao và khu vực với sai số ~1,5%.
- Cấu hình được quản lý thông qua một tệp YAML được tạo từ cấu hình sản phẩm, cho phép tải nóng (hot-loading) mà không cần triển khai lại.
- Tỷ lệ thu hồi trong phạm vi đạt đỉnh 86% nhưng ổn định quanh mức 64% trong những tháng khó khăn, cho thấy vẫn còn dư địa để cải thiện.
- Một kiểm tra sâu tổng hợp (synthetic deep-check) chạy mỗi 15 phút để xác thực toàn bộ đường ống cảnh báo từ khâu chèn dữ liệu đến khi đóng ticket.
Tại sao điều này quan trọng
Đối với các kỹ sư xây dựng nền tảng observability, nghiên cứu tình huống này minh họa rõ ràng các đánh đổi giữa tổng hợp kiểu batch và xử lý stream thời gian thực. Việc chuyển sang Flink cho phép nhóm tách rời chi phí khỏi số lượng tính năng được tích hợp, một yếu tố quan trọng đối với các sản phẩm SaaS đang phát triển. Việc sử dụng OpenTelemetry cho khả năng hiển thị end-to-end đảm bảo rằng bản thân hệ thống giám sát cũng có thể được giám sát, ngăn ngừa các lỗi âm thầm làm xói mòn niềm tin vào dashboard. Tuy nhiên, tỷ lệ thu hồi dao động nhắc nhở các nhà phát triển rằng dữ liệu nhanh hơn không tự động có nghĩa là logic phát hiện tốt hơn.
Lựa chọn kiến trúc sử dụng các sink idempotent và các topic Kafka chuyên dụng giải quyết các điểm đau phổ biến trong hệ thống phân tán: trùng lặp dữ liệu và tranh chấp tài nguyên. Bằng cách coi UID của operator như một API ổn định và điều chỉnh ranh giới autoscaler, nhóm đã loại bỏ các lần khởi động lại thường xuyên vốn trước đây làm gián đoạn trạng thái. Cách tiếp cận này cung cấp một khuôn mẫu cho các nhóm đang vật lộn với các ngăn xếp giám sát cũ, ồn ào, đắt đỏ và chậm chạp, vốn phụ thuộc vào hàng đợi dùng chung và bộ nhớ đệm trong bộ nhớ.
Bạn có thể làm gì
- Đánh giá đường ống sự kiện hiện tại của bạn để tìm các phụ thuộc vào hàng đợi dùng chung có thể gây ra vấn đề "hàng xóm ồn ào" trong giờ cao điểm.
- Cân nhắc sử dụng HyperLogLog hoặc các cấu trúc dữ liệu xác suất tương tự nếu bạn cần đếm số lượng riêng biệt trong các hệ thống phân tán.
- Triển khai lọc phía máy chủ ở cấp độ message bus để giảm khối lượng và chi phí xử lý downstream.
- Thiết kế các job xử lý stream với các sink idempotent và topic chuyên dụng để tránh trùng lặp dữ liệu và tranh chấp tài nguyên.


