0

Thiết kế pipeline dữ liệu crypto thời gian thực: độ mới, dự phòng và khả năng quan sát

Thiết kế pipeline dữ liệu thị trường crypto thời gian thực: độ mới, dự phòng và khả năng quan sát

Trong một sản phẩm dữ liệu crypto, “real-time” không chỉ có nghĩa là cập nhật nhanh. Một con số đến sớm nhưng không rõ nguồn, sai thời điểm hoặc bị lặp vẫn có thể gây hiểu nhầm. Khi xây dựng pipeline cho dữ liệu giá, khối lượng và biến động thị trường, tôi xem ba yếu tố là một thể thống nhất: độ mới của dữ liệu (freshness), khả năng chuyển nguồn (fallback) và khả năng quan sát hệ thống (observability).

Bài viết này trình bày một kiến trúc thực tế, độc lập với nhà cung cấp dữ liệu cụ thể.

1. Xác định “độ mới” trước khi tối ưu tốc độ

Mỗi bản ghi nên có ít nhất ba mốc thời gian:

  • event_time: thời điểm sự kiện được tạo tại nguồn;
  • received_at: thời điểm hệ thống nhận được dữ liệu;
  • published_at: thời điểm dữ liệu được phát tới client.

Từ đó có thể tính:

source_delay = received_at - event_time
processing_delay = published_at - received_at
end_to_end_delay = published_at - event_time

Nếu chỉ đo thời gian xử lý nội bộ, chúng ta có thể bỏ qua trường hợp nguồn đã gửi một bản ghi cũ. Với mỗi loại dữ liệu, cần định nghĩa một ngưỡng riêng. Ví dụ, ticker có thể bị coi là cũ sau vài giây, trong khi dữ liệu mô tả token có thể chấp nhận chu kỳ cập nhật dài hơn.

2. Chuẩn hóa dữ liệu ngay tại lớp adapter

Các nguồn thường khác nhau về tên cặp giao dịch, đơn vị, độ chính xác và định dạng timestamp. Vì vậy, mỗi nguồn nên có một adapter nhỏ chuyển dữ liệu về một schema chung:

type MarketTick = {
  symbol: string;
  price: string;
  volume24h?: string;
  eventTime: number;
  receivedAt: number;
  source: string;
  sequence?: string;
};

Giá và khối lượng nên được giữ dưới dạng decimal/string thay vì float nếu độ chính xác tài chính quan trọng. Adapter cũng là nơi phù hợp để kiểm tra schema, loại bỏ giá trị âm, timestamp nằm trong tương lai quá xa hoặc symbol không hợp lệ.

3. Khử trùng lặp và xử lý dữ liệu đến sai thứ tự

WebSocket có thể reconnect, REST polling có thể trả lại cửa sổ dữ liệu cũ, và message broker có thể giao lại một message. Do đó pipeline cần idempotency.

Khóa khử trùng lặp có thể là:

(source, symbol, sequence)

Nếu nguồn không cung cấp sequence, có thể dùng tổ hợp (source, symbol, event_time, price) trong một cửa sổ thời gian ngắn. Đồng thời, không nên giả định rằng mọi message đến đúng thứ tự. Mỗi symbol có thể giữ một watermark; bản ghi cũ hơn watermark chỉ được lưu cho mục đích audit, không ghi đè trạng thái hiện tại.

4. Fallback không đồng nghĩa với “chọn giá bất kỳ”

Một nguồn dự phòng cần được đánh giá theo trạng thái, không chỉ theo việc endpoint có trả HTTP 200 hay không. Health score có thể bao gồm:

  • tỷ lệ message hợp lệ;
  • độ trễ p50/p95;
  • thời gian kể từ bản cập nhật cuối;
  • số lần reconnect;
  • mức chênh lệch so với median của các nguồn còn lại.

Có thể mô hình hóa đơn giản:

score = w1 * freshness + w2 * availability + w3 * consistency

Khi chuyển nguồn, nên dùng hysteresis: chỉ chuyển sau N lần kiểm tra lỗi liên tiếp và chỉ quay lại nguồn chính khi nó ổn định trong một khoảng thời gian. Cách này tránh tình trạng hệ thống liên tục “nhảy” giữa hai nguồn.

Quan trọng hơn, client cần biết nguồn đã thay đổi. Trường sourcedata_status (live, degraded, stale) nên là một phần của response thay vì bị ẩn trong backend.

5. Đừng che giấu dữ liệu cũ

Khi không còn nguồn nào đạt ngưỡng, có ba lựa chọn:

  1. Trả lỗi và không hiển thị gì;
  2. Hiển thị giá trị gần nhất như thể nó vẫn mới;
  3. Hiển thị giá trị gần nhất kèm nhãn “stale” và thời điểm cập nhật.

Trong phần lớn giao diện thông tin, lựa chọn thứ ba minh bạch hơn. API có thể trả:

{
  "symbol": "BTC-USD",
  "price": "...",
  "source": "provider_b",
  "event_time": "...",
  "age_ms": 4200,
  "data_status": "degraded"
}

Giao diện nên hiển thị “cập nhật cách đây X giây” và cảnh báo rõ khi dữ liệu vượt ngưỡng. Không nên dùng dữ liệu stale để kích hoạt cảnh báo giá hoặc logic giao dịch tự động.

6. Observability theo symbol và theo nguồn

Dashboard tổng thể có thể trông bình thường dù một symbol hoặc một khu vực đang lỗi. Vì vậy metric nên có ít nhất hai chiều: sourcesymbol_group.

Những metric hữu ích:

  • market_data_age_ms;
  • invalid_message_total;
  • duplicate_message_total;
  • source_switch_total;
  • websocket_reconnect_total;
  • publish_latency_ms;
  • stale_response_total.

Log cần có correlation ID xuyên suốt từ adapter đến API. Với tracing, chỉ nên sample một phần traffic bình thường nhưng giữ tỷ lệ cao hơn cho request chậm hoặc dữ liệu bị đánh dấu degraded.

Alert cũng nên dựa trên triệu chứng người dùng nhìn thấy. Ví dụ, “95% symbol quan trọng có age_ms lớn hơn ngưỡng trong 3 phút” thường hữu ích hơn “CPU của worker vượt 80%”.

7. Thiết kế contract rõ ràng cho frontend

Frontend không nên tự suy đoán chất lượng dữ liệu chỉ từ timestamp. Backend nên cung cấp contract thống nhất:

type DataQuality = {
  status: "live" | "degraded" | "stale";
  source: string;
  eventTime: string;
  receivedAt: string;
  ageMs: number;
};

Nhờ đó web, mobile và API consumer có thể dùng cùng một quy tắc. Khi thay nhà cung cấp phía sau, hành vi của client vẫn nhất quán.

8. Kiểm thử các tình huống xấu trước khi production gặp chúng

Ngoài unit test cho adapter, nên có các bài test sau:

  • message trùng lặp;
  • timestamp đi lùi;
  • WebSocket ngắt và kết nối lại;
  • nguồn trả dữ liệu hợp lệ nhưng quá cũ;
  • hai nguồn lệch nhau bất thường;
  • rate limit hoặc timeout kéo dài;
  • chuyển nguồn rồi phục hồi nguồn chính;
  • toàn bộ nguồn mất kết nối.

Một replay harness từ dữ liệu đã ghi lại giúp tái hiện lỗi mà không phụ thuộc thị trường đang hoạt động thế nào tại thời điểm test.

Kết luận

Pipeline dữ liệu thị trường đáng tin cậy không được xây bằng một endpoint nhanh duy nhất. Nó cần schema chung, timestamp có ý nghĩa, idempotency, chính sách fallback có hysteresis, trạng thái chất lượng hiển thị cho người dùng và metric bám sát trải nghiệm thực tế.

Tôi là Mushegh Manukyan, Founder & CEO của ARMCP, làm việc tại Yerevan về dữ liệu thị trường, công cụ blockchain và khả năng tiếp cận thông tin crypto. Bài viết này tổng hợp các nguyên tắc kỹ thuật có thể áp dụng cho nhiều sản phẩm dữ liệu thời gian thực, không phụ thuộc vào một nhà cung cấp cụ thể.


All rights reserved

Viblo
Hãy đăng ký một tài khoản Viblo để nhận được nhiều bài viết thú vị hơn.
Đăng kí