Fetching Real-time OHLC
We now use the new Price Index Streams to fetch pre-aggregated OHLC data directly from Bitquery's GraphQL WebSocket API. This removes the need to manually calculate candlesticks from raw trade data.
To learn more about streaming data via graphQL, visit Bitquery subscriptions.
Imports and Configuration
import { createClient } from "graphql-ws";
import config from "./configs.json";
- createClient: From the
graphql-wslibrary, used to subscribe to GraphQL streams. - config: A local file storing your Bitquery API token.
let client;
/** Last emitted bar time and close — used to stitch new candles to the previous close. */
let lastEmittedBarTime = null;
let lastEmittedClose = null;
const BITQUERY_ENDPOINT =
"wss://streaming.bitquery.io/graphql?token=" + config.authtoken;
Subscription Query
The subscription below uses the Tokens cube, whose price blends every pool where the token is base. For a chart of one token, subscribe to Pairs with Ranking: { Position: { eq: 1 } } instead — the same Price.Ohlc fields, taken from the token's top market. Note that the top market can change during a stream, so read Market.Address from each message rather than assuming a fixed pool.
Run Stream on IDE We have used Solana as an example below, you can remove it and get data for all chains provided by the Price API
const subscriptionQuery = `
subscription{
Trading {
Tokens(
where: {Token: {Network: {is: "Solana"}, Address: {is: "6ft9XJZX7wYEH1aywspW5TiXDcshGc2W2SqBHN9SLAEJ"}}, Interval: {Time: {Duration: {eq: 1}}}}
) {
Block {
Time
}
Price {
Ohlc {
Open
High
Low
Close
}
}
Volume {
Base
Quote
}
Supply {
TotalSupply
MarketCap
FullyDilutedValuationUsd
}
}
}
}
`;
-
This query subscribes to pre-aggregated 1-second OHLC bars for a token on the Solana network (interval duration
1in thewhereclause). -
It requests:
- Block —
Time(bar timestamp) - Price.Ohlc —
Open,High,Low,Close - Volume —
Base,Quote - Supply —
TotalSupply,MarketCap,FullyDilutedValuationUsd
- Block —
Subscribing to the Stream
Keep the last emitted bar’s time and close in module-level variables. When the stream moves to a new candle (timestamp changes), set the new bar’s open to last candle’s close and widen high / low to include that price—same idea as historical bar continuity, but applied live as bars arrive.
export function subscribeToWebSocket(onRealtimeCallback) {
lastEmittedBarTime = null;
lastEmittedClose = null;
client = createClient({ url: BITQUERY_ENDPOINT });
const onNext = (data) => {
const tokenData = data.data?.Trading?.Tokens?.[0];
if (!tokenData) return;
const bar = {
time: new Date(tokenData.Block.Time).getTime(),
open: tokenData.Price.Ohlc.Open,
high: tokenData.Price.Ohlc.High,
low: tokenData.Price.Ohlc.Low,
close: tokenData.Price.Ohlc.Close,
volume: tokenData.Volume.Base,
};
const isNewCandle =
lastEmittedBarTime !== null && bar.time !== lastEmittedBarTime;
if (isNewCandle && lastEmittedClose != null) {
bar.open = lastEmittedClose;
bar.high = Math.max(bar.high, lastEmittedClose);
bar.low = Math.min(bar.low, lastEmittedClose);
}
lastEmittedBarTime = bar.time;
lastEmittedClose = bar.close;
onRealtimeCallback(bar);
};
client.subscribe(
{ query: subscriptionQuery },
{ next: onNext, error: console.error }
);
}
-
subscribeToWebSocket:
- Connects to Bitquery using
graphql-ws. - On each message, normalizes the bar for continuity when the interval rolls forward, then passes it to
onRealtimeCallback.
- Connects to Bitquery using
Unsubscribing from the Stream
export function unsubscribeFromWebSocket() {
if (client) {
client.dispose();
}
lastEmittedBarTime = null;
lastEmittedClose = null;
}
- unsubscribeFromWebSocket: Terminates the WebSocket connection and clears continuity state so a later reconnect does not stitch against stale closes.
Ready to run this in production?
Get an API key and run these queries in minutes, or talk to us about plans and enterprise delivery.