FDE PulseViệc làm FDE đang mở 316Mới đăng 7 ngày qua 10Chủ đề nổi bật: Đào tạo kỹ năng FDE tại Đông Nam Á

Tờ báo của nghề Forward Deployed Engineer

Bách khoa

Kafka và Flink cho FDE: ghép dòng dữ liệu của khách vào model chấm điểm thời gian thực

Model chạy tốt trên notebook là chuyện dễ. Khách sẽ hỏi bạn những câu khó hơn: có mất message không, có chấm trùng không, model chậm thì pipeline sẽ ra sao.

Đồ hoạĐường đi của một giao dịch tới điểm gian lận
  1. 1Kafka topic giao dịchKafka đưa sự kiện của khách vào Flink để xử lý
  2. 2KafkaSource + watermarkGắn thời điểm thật, chịu dữ liệu đến trễ, withIdleness cho partition vắng
  3. 3Async I/O gọi model serverNhiều request cùng lúc; capacity chạm trần thì tạo backpressure
  4. 4KafkaSink exactly-onceGhi trong transaction Kafka, commit khi checkpoint hoàn tất
  5. 5Topic điểm rủi roHệ thống duyệt giao dịch đọc điểm để ra quyết định

Mỗi chặng trả lời một câu hỏi của khách: thời gian, thông lượng, và chấm đúng một lần.

Đồ hoạ: FDE Times

Tóm tắt nhanh

  • Có hai cách đặt model: nhúng vào job Flink (không tốn lượt gọi mạng) hoặc gọi model server từ xa (chịu thêm độ trễ để được quản lý model tập trung).
  • Async I/O giúp một instance xử lý nhiều request cùng lúc, còn capacity giới hạn số request đang chờ và tạo backpressure khi chạm trần.
  • Khả năng khôi phục khi lỗi đến từ checkpoint chứ không đến từ offset đã commit; KafkaSink exactly-once commit transaction theo checkpoint.
Chia sẻLinkedInFacebookX

Thử hình dung khách của bạn là một cổng thanh toán. Đội data science giao cho bạn một model chấm điểm gian lận cho kết quả đẹp trên dữ liệu lịch sử, còn đội platform chỉ cho bạn một topic Kafka nơi giao dịch đang đổ về liên tục. Mỗi giao dịch phải có điểm rủi ro trước khi hệ thống duyệt kịp ra quyết định.

Ngân hàng và cổng thanh toán chính là nhóm điển hình dùng Flink gọi model từ xa để phân tích giao dịch theo thời gian thực, theo mô tả của Confluent. Trong kiến trúc đó, Kafka là lớp quen thuộc đưa dữ liệu vào Flink, còn Flink lo xử lý và gọi model.

Với một FDE, phần khó ít khi nằm ở model. Phần khó là ba câu khách sẽ hỏi ngay khi có sự cố đầu tiên: có mất message không, có giao dịch nào bị chấm hai lần không, model chậm thì chuyện gì xảy ra. Câu trả lời cho cả ba nằm trong vài dòng cấu hình mà bạn cần hiểu kỹ.

Model nên sống trong job hay ngoài job?

Kai Waehner, người viết nhiều về Kafka, chia cách triển khai thành hai kiểu. Kiểu embedded nhúng model thẳng vào ứng dụng xử lý luồng, nên mỗi lần chấm điểm không tốn một lượt gọi mạng. Kiểu remote gửi request tới model server qua RPC, API hoặc HTTP rồi chờ response.

Embedded Remote
Model chạy ở đâu Bên trong job Flink Trên model server riêng
Gọi mạng mỗi sự kiện Không Có, qua RPC, API hoặc HTTP
Quản lý model Đi cùng job: đổi model là deploy lại job Tập trung ở một chỗ
Cái giá phải trả Ít linh hoạt hơn Thêm độ trễ

Dòng “quản lý model” của cột embedded là suy luận từ chính cách nhúng: model nằm bên trong job, nên muốn thay phiên bản model thì bạn phải deploy lại cả job Flink.

Confluent mô tả remote inference là lựa chọn cho hệ thống thông lượng cao, chấp nhận độ trễ để đổi lấy sự linh hoạt. Ở site khách, nên hỏi trước tiên: ai sở hữu model và bao lâu họ retrain một lần?

Nếu đội data science muốn đưa phiên bản mới lên mà không phải chạm vào job Flink, hãy chọn remote. Phần còn lại của bài đi theo hướng đó.

Một pipeline chấm điểm gian lận, từng đoạn một

Đoạn đầu tiên là đọc dữ liệu và gắn thời gian cho nó.

env.enableCheckpointing(60_000);

KafkaSource<Txn> source = KafkaSource.<Txn>builder()
    .setBootstrapServers("broker:9092")
    .setTopics("payments.transactions")
    .setGroupId("fraud-scoring")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new TxnDeserializer())
    .setProperty("partition.discovery.interval.ms", "300000")
    .build();

WatermarkStrategy<Txn> wm = WatermarkStrategy
    .<Txn>forBoundedOutOfOrderness(Duration.ofSeconds(20))
    .withTimestampAssigner((txn, ts) -> txn.eventTimeMillis())
    .withIdleness(Duration.ofMinutes(1));

DataStream<Txn> txns = env.fromSource(source, wm, "payments");

Flink cần một WatermarkStrategy gồm hai phần: TimestampAssigner lấy thời điểm giao dịch thật từ dữ liệu, và WatermarkGenerator quyết định chờ dữ liệu đến trễ bao lâu. forBoundedOutOfOrderness với 20 giây nói rằng bạn chấp nhận giao dịch đến muộn tới chừng đó, một con số bạn nên thống nhất với khách thay vì tự đoán.

Dòng withIdleness cứu bạn khỏi một lỗi khó thấy. Nếu một partition Kafka không có dữ liệu, nó có thể kìm watermark của cả job đứng yên, và mọi phép tính theo thời gian sẽ chờ mãi. Đánh dấu input là idle sau một phút giải quyết đúng chuyện đó.

Dòng partition.discovery.interval.ms cũng đáng giải thích cho khách. KafkaSource tự phát hiện partition mới khi topic scale-out mà không cần restart job, và mặc định cứ 5 phút kiểm tra một lần, tức là đúng con số 300000 mili giây ở trên.

Đoạn thứ hai là gọi model. Cách gọi ở đây quyết định mỗi instance chấm được bao nhiêu giao dịch mỗi giây.

DataStream<Scored> scored = AsyncDataStream.unorderedWait(
    txns,
    new ScoreWithModelServer(),   // RichAsyncFunction dùng HTTP client bất đồng bộ
    500, TimeUnit.MILLISECONDS,   // timeout mỗi request
    100);                         // capacity: số request đang chờ tối đa

Toán tử Async I/O gửi request tới model server và nhận response theo kiểu non-blocking. Nhờ vậy một instance song song có thể xử lý nhiều request cùng lúc. Hãy làm một phép tính giả định: model server trả lời trong 50ms.

Gọi đồng bộ trong một map, mỗi instance chỉ chấm được 20 giao dịch mỗi giây. Với capacity 100, trần lý thuyết lên tới 2000, miễn là model server chịu nổi. Đừng biến con số 2000 thành lời hứa với khách: trần thật do sức chịu của model server và số instance song song của job quyết định, nên hãy đo trước khi cam kết.

Capacity còn có tác dụng thứ hai. Khi số request đang chờ chạm trần, Async I/O tạo backpressure ngược lên phía Kafka thay vì tích một backlog vô hạn trong bộ nhớ.

Đoạn cuối cùng là ghi kết quả.

KafkaSink<Scored> sink = KafkaSink.<Scored>builder()
    .setBootstrapServers("broker:9092")
    .setRecordSerializer(new ScoredSerializer("payments.fraud-scores"))
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("fraud-scoring-v1")
    .build();

scored.sinkTo(sink);

Ở chế độ EXACTLY_ONCE, KafkaSink ghi mọi message trong một transaction Kafka và chỉ commit khi checkpoint hoàn tất. Đó là lý do dòng enableCheckpointing ở đầu job không phải chuyện tùy chọn. transactionalIdPrefix phải là duy nhất, nên hãy đặt tên gắn với job và phiên bản.

Tự dựng ở site khách theo thứ tự nào?

Đừng ghép cả ba đoạn trong một buổi. Hãy bắt đầu bằng một job chỉ đọc topic và in ra, để chắc rằng deserializer đọc đúng dữ liệu và job có quyền truy cập topic trước khi thêm bất cứ thứ gì khác. Sau đó thêm watermark, nhìn vào Flink UI xem watermark có tiến lên trên mọi partition không.

Bước kế tiếp là thay model thật bằng một stub trả điểm cố định với độ trễ giả lập. Đo thông lượng ở vài mức capacity khác nhau và ghi lại, vì đó là bằng chứng để bạn thương lượng với đội vận hành model server về số request mỗi giây họ phải chịu. Chỉ khi các con số đã ổn mới nối model thật và bật exactly-once.

Những lỗi khiến pipeline hỏng mà không ai hay

Lỗi đầu tiên là nhìn offset đã commit trong Kafka rồi tin rằng đó là cơ chế khôi phục. KafkaSource có commit offset khi checkpoint hoàn tất, nhưng tài liệu Flink nói rõ nó không dựa vào offset đó để chịu lỗi. Offset đã commit chỉ để giám sát độ trễ của consumer; trạng thái thật nằm trong checkpoint.

Lỗi thứ hai là chạy hai job, chẳng hạn bản cũ và bản thử nghiệm, với cùng một transactionalIdPrefix. Yêu cầu prefix duy nhất có lý do, nên hãy đưa tên môi trường và phiên bản vào prefix ngay từ đầu.

Lỗi thứ ba là quên withIdleness trên một topic có partition lúc vắng lúc đông, ví dụ giao dịch ban đêm. Triệu chứng là các feature theo cửa sổ thời gian ngừng ra kết quả dù job vẫn “xanh”. Lỗi thứ tư là thấy backpressure rồi tăng capacity vô tội vạ, trong khi việc nên làm là hỏi xem model server có đang quá tải không.

Ghi kỹ năng này vào CV thế nào

Đừng viết “có kinh nghiệm Kafka, Flink”. Hãy viết một dòng có số đo từ chính bài tập của bạn: chuyển từ gọi model đồng bộ sang Async I/O, thông lượng tăng từ bao nhiêu lên bao nhiêu, exactly-once qua checkpoint. Kèm theo đó là một repo nhỏ có README trả lời đúng ba câu hỏi của khách: mất message không, chấm trùng không, model chậm thì sao.

Khách không mua một model, họ mua lời cam kết rằng mỗi giao dịch sẽ được chấm đúng một lần, kịp lúc. Nếu bạn chỉ ra được đúng dòng cấu hình đứng sau từng chữ trong lời cam kết đó, khách sẽ tin bạn.

5 nguồn
Đọc tiếp trên lộ trình · Chặng 2: Kỹ thuật rộngMLOps cho FDE: năm việc phải dựng xong trước khi model chạy trong hệ thống của kháchLần đầu khách gọi báo model dự báo sai, bạn cần trả lời được ngay ba câu: model nào đang chạy, sai từ bao giờ và lỗi nằm ở dữ liệu hay ở model.