Vnstock Logo

Khung streaming dữ liệu

Mở rộng

Mục lục

Ở thế hệ 5, vnstock_pipeline cung cấp một khung streaming mở (Streaming Framework), cho phép bạn xây dựng luồng nhận dữ liệu trong phiên giao dịch qua hai giao thức mạng WebSocket và Socket.IO. Dữ liệu cập nhật theo độ trễ của nguồn bạn kết nối.

Khung kết nối mở (Bring Your Own Source)

Thư viện không đóng gói sẵn nguồn dữ liệu nào cho luồng trong phiên. Thư viện lo phần quản lý kết nối, bộ đệm (buffering), định dạng cột, làm phẳng dữ liệu (flattener) và ghi đĩa. Bạn tự cung cấp địa chỉ kết nối (URI) và bộ giải mã cho nguồn mà bạn có quyền truy cập, và tuân theo điều kiện của nguồn đó.

Các thành phần trong khung streaming

Kiến trúc streaming của vnstock_pipeline được chia thành các lớp chức năng rõ ràng:

  • Client kết nối:
    • SocketIOClient: Quản lý bắt tay (handshake), heartbeat (ping/pong) và duy trì phiên kết nối chuẩn Socket.IO (EIO=4...).
    • BaseWebSocketClient: Lớp cơ sở dùng cho các kết nối WebSocket thuần tuý.
  • Bộ phân tích dữ liệu (BaseDataParser): Lớp trừu tượng cho phép bạn chuyển đổi các thông điệp thô (raw text/JSON) từ nguồn thành cấu trúc từ điển Python đã chuẩn hoá tên trường.
  • Bộ làm phẳng và căn chỉnh cột (DataFlattener): Đưa các trường lồng nhau trong JSON về dạng phẳng một cấp và sắp xếp theo danh sách cột cố định.
  • Bộ xử lý thành phẩm (CSVProcessor): Ghi luồng dữ liệu nhận được trong phiên trực tiếp vào tệp CSV trên đĩa với dung lượng bộ đệm tối ưu.
  • Quản lý phiên (SessionManager): Nhận diện khung giờ giao dịch của thị trường chứng khoán Việt Nam, duy trì kết nối trong phiên và khôi phục kết nối khi mạng bị gián đoạn.

Hướng dẫn kết nối và nhận dữ liệu

Dưới đây là ví dụ hoàn chỉnh về cách thiết lập một luồng truy xuất dữ liệu bằng SocketIOClient:

Python
import asyncio
from vnstock_pipeline.stream import BaseDataParser, CSVProcessor, SocketIOClient

# Bước 1: Định nghĩa bộ phân tích dữ liệu phù hợp với nguồn của bạn
class CustomMarketParser(BaseDataParser):
    def parse_data(self, raw_data: dict) -> dict:
        """Trích xuất và chuẩn hoá các trường dữ liệu cần thiết."""
        # Giả sử cấu trúc JSON nhận về có dạng: {"event": "price_tick", "data": {...}}
        tick_data = raw_data.get("data", {})
        return {
            "symbol": tick_data.get("symbol"),
            "price": tick_data.get("last_price"),
            "volume": tick_data.get("last_vol"),
            "time": tick_data.get("time"),
            "event_type": raw_data.get("event", "tick")
        }

async def main():
    # Bước 2: Khởi tạo client với địa chỉ nguồn cấp của bạn
    source_uri = "wss://your-streaming-endpoint.example.com/socket.io/?EIO=4&transport=websocket"
    client = SocketIOClient(
        uri=source_uri,
        parser=CustomMarketParser()
    )

    # Bước 3: Gửi thông điệp đăng ký nhận dữ liệu cho các mã quan tâm
    subscribe_msg = '42["subscribe",{"symbols":["ACB","HPG","VNM"]}]'
    client.add_raw_message(subscribe_msg)

    # Bước 4: Gắn bộ xử lý ghi dữ liệu vào tệp CSV
    client.add_processor(CSVProcessor(output_dir="./data/streaming_feed"))

    # Bước 5: Kết nối và khởi động tiến trình giám sát phiên
    await client.connect()
    await client.start_session_monitoring()

# Khởi chạy luồng sự kiện bất đồng bộ
# asyncio.run(main())

Tuỳ biến xử lý luồng dữ liệu

Ngoài việc ghi ra CSV bằng CSVProcessor, bạn có thể tự viết bộ xử lý tuỳ biến để chuyển dữ liệu vào cơ sở dữ liệu hàng đợi (Queue), hệ thống cảnh báo (Telegram/Discord) hoặc giao diện hiển thị:

Python
from vnstock_pipeline.stream.processors import BaseProcessor

class TelegramAlertProcessor(BaseProcessor):
    def process_data(self, data: dict):
        # Kiểm tra điều kiện và gửi cảnh báo khi có giao dịch lớn
        volume = data.get("volume", 0)
        if volume > 100000:
            print(f"Lệnh lớn: {data.get('symbol')} khớp {volume} cổ phiếu tại giá {data.get('price')}")

# Gắn thêm bộ xử lý cảnh báo vào client
client.add_processor(TelegramAlertProcessor())

Khả năng chịu lỗi và quản lý kết nối

  • Cơ chế Heartbeat: Định kỳ gửi gói tin kiểm tra kết nối với máy chủ nguồn để tránh bị đóng ngắt do hết thời gian chờ (timeout).
  • Khôi phục kết nối: Khi mạng Internet gặp sự cố hoặc máy chủ nguồn ngắt kết nối đột ngột, client sẽ áp dụng thuật toán thử lại có độ trễ luỹ tiến (exponential backoff) để thiết lập lại liên kết.
  • Tiết kiệm tài nguyên: Sau giờ đóng cửa thị trường (15:00), SessionManager sẽ đưa kết nối vào trạng thái chờ nghỉ nhằm tiết kiệm băng thông và tài nguyên máy tính của bạn.