You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RXJS5实现流排队至事件触发后正常处理及GDAX订单簿同步问询

Alright, let's walk through how to solve this exact problem with RxJS 5—this is a super common scenario when building order books from exchange APIs, so I’ve got you covered.

RxJS 5 Solution: Queue WebSocket Events Until Snapshot Sync Completes

Core Idea

Here's the breakdown of what we need to accomplish:

  • First fetch and process the order book snapshot (Source A) to get its latest sequence number
  • Buffer all incoming WebSocket events (Source B) while the snapshot is being processed
  • Once the snapshot is done and we emit the SYNCED signal, filter out any buffered events older than the snapshot's sequence number, then replay the valid ones
  • After sync, let new WebSocket events pass through normally without buffering

Step-by-Step Implementation

1. Define Your Data Sources

First, let's set up references to your actual data streams (I'll include mock implementations that match GDAX's structure):

// Source A: Order book snapshot API call (returns snapshot with sequence number)
const sourceA$ = Rx.Observable.defer(() => 
  fetch('https://api.gdax.com/products/BTC-USD/book?level=3')
    .then(res => res.json())
).shareReplay(1); // Share the snapshot so we don't fetch it twice

// Source B: GDAX WebSocket real-time events
const sourceB$ = Rx.Observable.webSocket('wss://ws-feed.gdax.com')
  .map(message => message); // Adjust mapping to match your message parsing needs

// SYNCED signal: Emitted once the snapshot is fully processed
const synced$ = sourceA$
  .do(snapshot => {
    // Build your initial order book here using the snapshot data
    console.log('Snapshot processed - latest sequence:', snapshot.sequence);
  })
  .mapTo('SYNCED'); // Send the sync signal after processing

2. Buffer and Filter Source B Events

We'll use bufferWhen to hold onto Source B events until the synced$ signal fires. Then we'll filter out outdated events and merge them with the live stream post-sync:

// Buffer Source B events until the SYNCED signal is received
const bufferedB$ = sourceB$.bufferWhen(() => synced$.take(1));

// Filter the buffered events to keep only those newer than the snapshot
const filteredBufferedB$ = bufferedB$
  .mergeMap(bufferedEvents => {
    return sourceA$.map(snapshot => {
      return bufferedEvents.filter(event => 
        event.sequence > snapshot.sequence // Discard events older than the snapshot
      );
    });
  })
  .mergeAll(); // Flatten the filtered array into individual events

// Combine filtered buffered events with live post-sync events
const processedB$ = Rx.Observable.merge(
  filteredBufferedB$,
  sourceB$.skipUntil(synced$) // Let live events flow after sync
);

3. Assemble the Full Workflow

Finally, we'll chain everything together to ensure the order book is built correctly:

// Main workflow: First process the snapshot, then handle filtered/live events
const orderBookWorkflow$ = Rx.Observable.concat(
  sourceA$.do(snapshot => /* Initialize your order book here */),
  processedB$.do(event => /* Update your order book with each event */)
);

// Start the workflow by subscribing
orderBookWorkflow$.subscribe(
  () => {},
  err => console.error('Order book error:', err),
  () => console.log('Workflow completed')
);

Key Tips for Production

  • Error Handling: Add catch() to sourceA$ to retry snapshot fetches if they fail—you don't want a single API error to break the entire sync.
  • WebSocket Connection: Ensure your WebSocket is connected before fetching the snapshot. You can add a check for the connection status or use sourceB$.take(1).delay(0) to wait for the initial connection handshake.
  • Memory Management: If your buffer gets large (e.g., slow snapshot fetch), consider adding a max buffer size to avoid memory issues—though in practice, GDAX snapshots are fast enough that this rarely matters.

内容的提问来源于stack exchange,提问作者kyle

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 06:57:19