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
SYNCEDsignal, 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()tosourceA$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
相关产品推荐
相关产品推荐

