首页 / 文章 / 基于SSE、虚拟化技术及20个市场预算的响应式交易用户界面

基于SSE、虚拟化技术及20个市场预算的响应式交易用户界面

协调React虚拟化行、中央SSE订阅规划器以及外部报价存储,从而使未处理的订单能够实时获取价格,同时不会超出连接限制。

1870 词

交易记录表可以列出数千种金融工具,但屏幕上仅显示几行数据。行情报价会持续更新,用户可滚动查看,而当相关市场离开可视区域后,未平仓订单仍可能需要实时价格。

这些需求存在时间上的差异:

Scrolling changes visibility.
Subscription policy changes network demand.
SSE changes quote data.
React changes what is displayed.

为每个需求指定唯一的负责人。下文的示例仅展示了这种划分方式,属于架构层面的拆分,并非完整的交易客户端实现。

在选择实现方案之前先明确规则

假设后端允许在单个SSE连接上最多使用20个市场ID。那么相应的策略如下:

  • 每个浏览器标签页对应一个活跃的EventSource
  • 该连接上使用的不同市场ID数量不得超过20个
  • 未平仓订单的优先级高于普通列表行
  • 需求变化时无需取消消费者注册
  • 实时报价独立于 React 组件状态存在
  • 传入的帧在修改存储数据前会经过验证
  • 提交时重新读取最新的报价,而不仅仅是最后显示的那一条
  • 已过期或断开的报价不可用于提交
  • 架构设计

    Virtualized list ─── visible and nearby market IDs ──┐
                                                       ├─► Subscription manager
    Open order ticket ─── selected market ID ───────────┘             │
                                                                     ▼
                                                              One SSE connection
                                                                     │
                                                              Validated messages
                                                                     ▼
                                                                 Quote store
                                                                /           \
                                                               ▼             ▼
                                                      React snapshots   Order validation
    

    React 组件负责声明需求并读取快照,它们不负责打开或管理套接字。

    1. 独立于 React 管理当前报价

    经过验证的报价数据格式如下:

    type Quote = Readonly<{
      marketId: string;
      bid: string;       // Decimal strings preserve wire precision.
      ask: string;
      version: number;
      staleAt: number;   // Server expiry timestamp in milliseconds.
      tradable: boolean;
    }>;
    

    快照会同时暴露数据是否就绪以及具体内容:

    type QuoteSnapshot =
      | { state: "waiting"; quote: null }
      | { state: "fresh" | "stale"; quote: Quote };
    

    存储层的结构保持简洁:

    interface QuoteStore {
      receive(quote: Quote): void;
    
      // Immediate state, used by commands.
      readLatest(marketId: string): QuoteSnapshot;  // Cached visual snapshots, used by React.
      select(marketId: string): {
        subscribe(notify: () => void): () => void;
        getSnapshot(): QuoteSnapshot;
      };  invalidate(marketIds: readonly string[]): void;
      retain(marketIds: readonly string[]): void;
    }
    

    接收逻辑会在接受更新前比较版本信息:

    function receiveQuote(next: Quote) {
      const previous = latestQuotes.get(next.marketId);
    
      // This example assumes strictly increasing quote versions.
      if (previous && next.version <= previous.version) return;  latestQuotes.set(next.marketId, next);  // Independent of whether another price message arrives.
      scheduleQuoteExpiry(next.marketId, next.staleAt);  // Ordinary prices can wait for the next visual publication.
      scheduleVisualPublication();
    }
    

    到期计时器、失效处理以及向订阅者的信息推送功能均位于存储系统中。历史记录并非无限保存——仅保留每个市场的最新报价:

    Market A, version 10
            ↓ replaced
    Market A, version 11
            ↓ replaced
    Market A, version 12
    

    如果服务器在不同类型的帧之间重复使用版本空间,还原器必须根据相应规则进行合并。

    2. 中央统一分配20个订阅资格

    消费者会按优先级声明需求:

    type Demand = {
      marketIds: readonly string[];
      priority: number;
    };
    
    const PRIORITY = {
      orderTicket: 0,
      visible: 1,
      nearby: 2,
    } as const;const MAX_MARKETS = 20;
    

    规划器会消除重复项,并保留每个市场中优先级最高的那个需求:

    function planSubscriptions(
      demands: readonly Demand[],
      activeIds: readonly string[],
    ): string[] {
      const priorities = new Map<string, number>();
    
      for (const demand of demands) {
        for (const marketId of demand.marketIds) {
          priorities.set(
            marketId,
            Math.min(
              priorities.get(marketId) ?? Infinity,
              demand.priority,
            ),
          );
        }
      }  if (priorities.size === 0) return [];  const active = new Set(activeIds);  const requested = [...priorities.keys()].sort(
        (left, right) =>
          priorities.get(left)! - priorities.get(right)! ||
          Number(active.has(right)) - Number(active.has(left)),
      );  return [...new Set([...requested, ...activeIds])]
        .slice(0, MAX_MARKETS)
        .sort();
    }
    

    当需求量超过20个市场时,多余的记录会以可见状态保持待处理状态。用户界面不得伪装这些记录已获得实时覆盖。

    3. 在不重启React效应的情况下更新需求

    每个消费者都会有一个固定的负责人:

    function useMarketDemand(
      marketIds: readonly string[],
      priority: number,
    ) {
      const [owner] = useState(() =>
        subscriptionManager.createOwner(),
      );
    
      const key = JSON.stringify([...new Set(marketIds)]);
      const stableIds = useMemo<string[]>(
        () => JSON.parse(key),
        [key],
      );  useEffect(() => {
        owner.update({
          marketIds: stableIds,
          priority,
        });
      }, [owner, stableIds, priority]);  useEffect(() => {
        return () => owner.dispose();
      }, [owner]);
    }
    

    ID变更会直接更新该负责人;组件卸载时则会清理相关记录。

    快速变化的需求会在固定的调度时间窗口内被批量处理:

    let reconciliationTimer:
      ReturnType<typeof setTimeout> | undefined;
    
    function scheduleReconciliation() {
      if (reconciliationTimer !== undefined) return;  reconciliationTimer = setTimeout(() => {
        reconciliationTimer = undefined;
        reconcileSubscriptions();
      }, 150);
    }
    

    由于该时间窗口是固定的,而非可滑动的延迟机制,因此无法通过持续滚动来无限推迟对账操作。首次连接和最终断开仍可能立即执行。

    4. 安全地替换 SSE 连接

    新的查询参数意味着需要一个新的 EventSource。在打开新流之前,必须先关闭之前的流:

    let source: EventSource | undefined;
    let generation = 0;
    let clearWatchdog: (() => void) | undefined;
    
    function replaceConnection(marketIds: string[]) {
      const currentGeneration = ++generation;  clearWatchdog?.();
      source?.close();
      source = undefined;  quoteStore.retain(marketIds);
      quoteStore.invalidate(marketIds);  if (marketIds.length === 0 || !navigator.onLine) return;  const query = new URLSearchParams();  for (const marketId of marketIds) {
        query.append("marketId", marketId);
      }  const nextSource = new EventSource(
        `/api/v1/stream?${query.toString()}`,
      );  source = nextSource;
      const permittedIds = new Set(marketIds);  const watchdog = createSilenceWatchdog(() => {
        if (currentGeneration !== generation) return;    nextSource.close();
        reconcileSubscriptions({ force: true });
      });  clearWatchdog = watchdog.stop;
      watchdog.reset();  nextSource.addEventListener("quote", event => {
        if (currentGeneration !== generation) return;    // Parses JSON and validates it against the feed contract.
        const quote = parseQuoteMessage(event);    if (!quote || !permittedIds.has(quote.marketId)) {
          quoteStore.invalidate(marketIds);
          return;
        }    watchdog.reset();
        quoteStore.receive(quote);
      });  nextSource.addEventListener("heartbeat", () => {
        if (currentGeneration === generation) {
          watchdog.reset();
        }
      });  nextSource.onerror = () => {
        if (currentGeneration !== generation) return;    quoteStore.invalidate(marketIds);    if (nextSource.readyState === EventSource.CLOSED) {
          watchdog.stop();
          reportConnectionError();
        }    // Recoverable failures are retried by native EventSource.
      };
    }
    

    看门狗机制、解析器以及连接状态检测功能都位于独立的模块中。版本计数器会将来自过时套接字的延迟事件剔除。原生 EventSource 会在可恢复的错误情况下自动重新连接;而调用 .close() 则会终止该实例。有关浏览器行为的更多信息,请参阅 MDN 关于服务器推送事件的指南。

    在线/离线状态以及页面生命周期监听器都位于管理器层级:当设备离线时,会取消报价并关闭套接字;恢复连接后会根据当前需求重新建立连接。新的连接会在允许提交之前等待新的报价生成。

    5. 对列表进行虚拟化并呈现其视口

    虚拟化功能可限制已挂载的DOM元素数量;而订阅规划器则负责单独控制流式数据的市场数据。通过TanStack Virtual,可见行可以请求比周围超出显示范围的行更高的优先级:

    function MarketList({
      markets,
      onOpenTicket,
    }: {
      markets: Array<{ marketId: string; name: string }>;
      onOpenTicket(marketId: string, side: "BUY" | "SELL"): void;
    }) {
      const [viewport, setViewport] =
        useState<HTMLDivElement | null>(null);
    
      const virtualizer = useVirtualizer({
        count: markets.length,
        getScrollElement: () => viewport,
        getItemKey: index => markets[index]!.marketId,
        estimateSize: () => 72,
        overscan: 2,
      });  const rows = virtualizer.getVirtualItems();
      const range = virtualizer.range;  const visibleIds = range
        ? markets
            .slice(range.startIndex, range.endIndex + 1)
            .map(market => market.marketId)
        : [];  const nearbyIds = rows.map(
        row => markets[row.index]!.marketId,
      );  useMarketDemand(visibleIds, PRIORITY.visible);
      useMarketDemand(nearbyIds, PRIORITY.nearby);  return (
        <div
          ref={setViewport}
          role="region"
          aria-label="Markets"
          tabIndex={0}
          style={{ height: 560, overflow: "auto" }}
        >
          <div
            style={{
              height: virtualizer.getTotalSize(),
              position: "relative",
            }}
          >
            {rows.map(row => (
              <div
                key={row.key}
                style={{
                  position: "absolute",
                  top: 0,
                  left: 0,
                  width: "100%",
                  height: 72,
                  transform: `translateY(${row.start}px)`,
                }}
              >
                <MarketRow
                  market={markets[row.index]!}
                  onOpenTicket={onOpenTicket}
                />
              </div>
            ))}
          </div>
        </div>
      );
    }
    

    这些示例假设行高是固定的;若行高可变,则需要先进行测量。虚拟化工具负责挂载各行,而订阅策略仍由应用程序自行管理(详见TanStack Virtual的React文档)。当行离开常规渲染窗口时,键盘焦点仍应确保该行保持可聚焦状态。

    6. 通过各自的快照来渲染每个市场

    function useQuote(marketId: string) {
      const selection = useMemo(
        () => quoteStore.select(marketId),
        [marketId],
      );
    
      return useSyncExternalStore(
        selection.subscribe,
        selection.getSnapshot,
      );
    }
    

    在那个市场的状况发生实质性变化之前,select会一直返回相同的快照引用——这正是React对与useSyncExternalStore一起使用的外部存储所期望的行为。

    const MarketRow = memo(function MarketRow({
      market,
      onOpenTicket,
    }: {
      market: { marketId: string; name: string };
      onOpenTicket(marketId: string, side: "BUY" | "SELL"): void;
    }) {
      const snapshot = useQuote(market.marketId);
      const quote = snapshot.quote;
    
      const available =
        snapshot.state === "fresh" && quote?.tradable;  return (
        <article>
          <strong>{market.name}</strong>      <button
            disabled={!available}
            onClick={() => onOpenTicket(market.marketId, "SELL")}
          >
            Sell · {quote?.bid ?? "—"}
          </button>      <button
            disabled={!available}
            onClick={() => onOpenTicket(market.marketId, "BUY")}
          >
            Buy · {quote?.ask ?? "—"}
          </button>      {snapshot.state !== "fresh" && (
            <span>Prices updating</span>
          )}
        </article>
      );
    });
    

    市场A的报价变动只会通知市场A内的读取器,而不会替换父列表中的整个市场数组。按钮仅用于提交工单,并不会默默地完成订单。

    该工单会明确说明自身的需求:

    function OrderTicket({ marketId }: { marketId: string }) {
      useMarketDemand([marketId], PRIORITY.orderTicket);
    
      const snapshot = useQuote(marketId);  // Render quantity, side, current quote, review and submit controls.
      // ...
    }
    

    滚动网格并不会取消工单中的需求。仅删除该工单只会移除对应的负责人,而可见的行或附近的额外扫描条目可能仍需要相同的金融工具。

    在买卖操作前务必立即重新查看实时报价

    上次绘制的行内容可能会使权威存储稍有延迟。因此,提交路径需要重新读取:

    async function submitOrder(intent: {
      clientOrderId: string; // Stable for retries of this exact intent.
      marketId: string;
      side: "BUY" | "SELL";
      quantity: string;
      reviewedVersion: number;
    }) {
      const snapshot = quoteStore.readLatest(intent.marketId);
      const quote = snapshot.quote;
    
      if (
        !navigator.onLine ||
        snapshot.state !== "fresh" ||
        !quote ||
        !quote.tradable ||
        Date.now() >= quote.staleAt
      ) {
        throw new Error("Wait for a fresh, tradable quote.");
      }  if (quote.version !== intent.reviewedVersion) {
        throw new Error("The quote changed. Review it again.");
      }  const limitPrice =
        intent.side === "BUY" ? quote.ask : quote.bid;  const response = await fetch("/api/v1/orders", {
        method: "POST",
        headers: { "Content-Type": "application/json" },
        body: JSON.stringify({
          ...intent,
          type: "LIMIT",
          limitPrice,
        }),
      });  return parseOrderResponse(response);
    }
    

    会对响应进行解析,以判断其为被接受、被拒绝还是存在歧义。若要重新尝试同一操作,则会使用其原有的客户端订单标识符。授权、报价验证、接受以及成交决策均由服务器端决定。价格限制仅能约束允许的执行操作——无法保证一定能够成交。

    两个时间决策

    决策项 目的
    约150毫秒的订阅窗口 在更换连接之前,先处理快速的视图切换情况
    约100毫秒的视觉更新间隔 限制普通价格信息的重新渲染频率

    传入的消息仍会立即更新当前存储。可用性状态的切换无需经历常规的视觉延迟。150毫秒的时间窗口并非为了减缓滚动、点击操作,或是已经通过活跃连接传来的价格更新速度。

    内存与清理

    目录缓存(如React Query)需要独立的保留规则。虚拟化列表虽然仅显示十行数据,但仍会缓存数千条已加载的记录。需将各功能模块区分开来:虚拟化负责渲染界面,调度器管理订阅容量,存储层掌控数据状态,而React则负责绘制重要的快照。