返回
RSS ClickHouse Blog AI 逐段翻译 发布 2026-09-04 08:00 收录于 09-07

用ClickHouse和Massive构建实时市场数据应用

DataHot 速览

本文介绍如何使用ClickHouse与Massive(原Polygon.io)构建可扩展至每秒数千事件的实时金融行情应用。文中解释了tick、quote、trade等概念,并演示了如何接入行情流、写入ClickHouse、利用物化视图预聚合,以及用Node.js和React实现实时可视化。

为什么值得关注:实时行情是典型的高频流分析场景,本文给出了从数据接入、存储到可视化的完整技术路径,对数据从业者构建低延迟分析应用有参考价值。

本文目录 12 节
  1. 什么是tick?
  2. 访问实时市场数据
  3. 将数据摄取到ClickHouse
  4. 在ClickHouse中建模Tick数据
  5. 摄取策略
  6. 可视化实时市场数据
  7. 扩展和实用技巧
  8. 扩展摄取
  9. 监控摄取延迟
  10. 利用物化视图
  11. 运行完整示例
  12. 结论

译文

AI 逐段翻译

实时逐笔数据应用是实时分析的一个经典例子。就像在Web应用中跟踪用户行为或监控物联网设备的指标一样,它们涉及需要以低延迟摄取、存储和查询的高频事件流。

在金融市场中,区别在于紧迫性。即使几秒钟的延迟也可能将盈利的交易变成亏损。每一笔交易和报价更新都会生成一个新的tick,并且跨多个品种每秒可能达到数千个。

ClickHouse非常适合这种类型的工作负载。它处理高频插入、基于时间的查询和低延迟查询。内置压缩有助于减少存储开销,即使每个品种有数十亿行。物化视图可用于在写入时预聚合或重新组织数据,优化查询性能,而无需单独的处理层。

在本文中,我们将演示如何使用Massive(原Polygon.io)访问市场数据,并使用ClickHouse实时存储和查询tick。我们将结合Node.js进行后端操作,React用于实时可视化。让我们开始吧。

什么是tick?

在开始之前,理解报价和交易的含义是有帮助的。报价代表市场参与者愿意买入或卖出证券的当前价格。具体而言,它包括最佳买价(有人愿意支付的最高价格)和最佳卖价(有人愿意卖出的最低价格)。随着新订单进入或离开市场,这些价格会持续更新。

symbxbpbsaxapascitqzinserted_at
SPY12602.73211602.7461[1,93]174958347820663152225NYSE1749583479396

另一方面,交易是买方和卖方之间的实际交易。当有人同意当前的卖价或买价,订单被匹配并执行时,就会发生交易。交易记录包括执行价格、交易量和时间戳。

symixpsctqztrfitrftinserted_at
SPY5298352503482511607.261[12,37]175084225597222126NYSE001750842257842

Tick数据通常来自两个流。一个包含报价更新,另一个包含交易执行。两者对于理解市场行为都至关重要,但在分析和策略开发中用途不同。

访问实时市场数据

现在我们知道要摄取的数据类型,让我们看看如何访问它。我们首先需要找到并订阅一个股票市场API。有很多可用的,我们选择用于构建此演示的是Massive。选择一个明确包含股票tradesquotes通过WebSockets的计划。仅有REST API密钥或仅包含聚合数据的计划是不够的。在订阅前检查当前渠道权限。

WebSockets对于流式市场数据至关重要,因为它们消除了轮询REST API的延迟和开销。无需为每个数据请求建立新连接,也不会在调用之间错过tick,WebSockets保持持久连接,在数据可用时立即推送,这对于毫秒级至关重要的高频市场数据非常关键。

开始使用MassiveAPI摄取数据相当直接,只需与/stocks端点建立连接,使用您的Massive API密钥进行身份验证,然后开始处理消息。

以下独立的Node.js诊断程序进行身份验证,在身份验证成功后订阅,并记录收到的事件。在安装依赖并配置.env后,从示例目录运行它。在启动完整应用程序前停止它;账户连接限制可能阻止多个客户端使用同一数据源。

1require("dotenv").config();2constWebSocket = require("ws");34if (!process.env.MASSIVE_API_KEY) thrownewError("Set MASSIVE_API_KEY in .env");5const symbols = (process.env.MASSIVE_SYMBOLS || "AAPL,MSFT,NVDA")6  .split(",").map((symbol) => symbol.trim()).filter(Boolean);7const ws = newWebSocket("wss://socket.massive.com/stocks");89ws.on("open", () => {10  ws.send(JSON.stringify({ action: "auth", params: process.env.MASSIVE_API_KEY }));11});12ws.on("message", (data) => {13const payload = JSON.parse(data.toString());14for (const row of payload) {15if (row.ev === "status") {16console.log(row.status, row.message);17if (row.status === "auth_success") {18        ws.send(JSON.stringify({19action: "subscribe",20params: symbols.flatMap((symbol) => [`T.${symbol}`, `Q.${symbol}`]).join(","),21        }));22      } elseif (row.status === "auth_failed") {23        ws.close();24      }25    } elseif (row.ev === "T" || row.ev === "Q") {26console.log(row);27    }28  }29});30ws.on("error", (error) =>console.error(error.message));31process.on("SIGINT", () => ws.close());

将数据摄取到ClickHouse

在ClickHouse中建模Tick数据

Tick数据建模相对直接,因为它只包含两种事件类型,一种用于trade,另一种用于quote,每种类型都有少量大多为数字的字段。以下是创建两个独立表的DDL。t是毫秒级的SIP时间戳,而q是每个符号的序列号。条件和指标数组使用UInt32,因为值可能超过255。交易模式还保留可选的ds(小数数量,作为文本)和pt(参与者时间戳)字段。缺少的可选字段使用默认值。当前报价数量以股表示。示例使用Float64表示价格,并对整数交易量s求和;在您的应用程序需要的地方,请使用合适的小数类型和分数数量处理。

1CREATE TABLE IF NOTEXISTS quotes2(3    `sym` LowCardinality(String),4    `bx` UInt8,5    `bp` Float64,6    `bs` UInt64,7    `ax` UInt8,8    `ap` Float64,9    `as` UInt64,10    `c` UInt8,11    `i` Array(UInt32),12    `t` UInt64,13    `q` UInt64,14    `z` Enum8('NYSE'=1, 'AMEX'=2, 'Nasdaq'=3),15    `inserted_at` UInt64 DEFAULT toUnixTimestamp64Milli(now64())16)17ENGINE = MergeTree18ORDERBY (sym, t - (t %60000));1920CREATE TABLE IF NOTEXISTS trades21(22    `sym` LowCardinality(String),23    `i` String,24    `x` UInt8,25    `p` Float64,26    `s` UInt64,27    `ds` String DEFAULT'',28    `pt` UInt64 DEFAULT0,29    `c` Array(UInt32),30    `t` UInt64,31    `q` UInt64,32    `z` Enum8('NYSE'=1, 'AMEX'=2, 'Nasdaq'=3),33    `trfi` UInt64,34    `trft` UInt64,35    `inserted_at` UInt64 DEFAULT toUnixTimestamp64Milli(now64())36)37ENGINE = MergeTree38ORDERBY (sym, t - (t %60000));

数据量可能迅速增长。

开发时订阅一个小型符号列表,然后测量吞吐量后再订阅整个市场。选择有效的排序键对于性能至关重要。在本例中,行首先按sym(股票代码)排序,将同一符号的所有事件分组在一起。在每个符号组内,行按有效的排序键对性能至关重要。在这种情况下,行首先按 sym(股票代码)排序,将同一代码的所有事件分组。在每个代码组内,行按t - (t % 60000)排序,这创建了1分钟的时间桶。这种方法在我们的案例中效果很好,因为我们按符号聚合数据来生成可视化。分钟桶将附近的事件分组,而查询仍然使用原始时间戳和序列号来确定事件顺序。根据您的工作负载对排序键和过滤器进行基准测试;仅凭时间桶并不能保证每次查询都高效剪枝。

摄取策略

有几种方法可以为此类应用程序设计摄取管道,包括使用消息队列如Kafka。然而,为了最小化延迟,通常最好保持系统简单,并尽可能将数据直接从WebSocket连接推送到ClickHouse。

一旦该设置就绪,下一步是选择正确的摄取方法。ClickHouse支持同步和异步插入。

对于同步摄取,数据在客户端分批发送。批次大小应在内存使用、延迟和系统开销之间取得平衡。较大的批次会减少插入请求次数并提高吞吐量,但可能增加内存使用并延迟单条记录的处理。较小的批次可减少内存压力,但可能因生成过多小型数据分区而给ClickHouse带来更多负载。

使用异步摄取时,数据会持续发送到ClickHouse,批处理在内部处理。传入的记录首先写入内存缓冲区,然后根据可配置阈值刷新到存储。当客户端批处理不切实际时(例如数据来自许多小型客户端),此方法很有用。

在我们的案例中,同步摄取更合适。由于只有一个客户端从WebSocket API推送数据,因此可以在客户端管理批处理,以便更好地控制性能和资源使用。

这个紧凑的批处理类展示了摄取路径。创建一个TickBatcher,并将传入的WebSocket帧传递给其handleMessage方法,并结合上述认证处理器。在等待close()之前停止传入的消息。完整示例增加了有界并发插入、重连处理、关闭期限和指标。两者都使用同步插入,因此成功意味着ClickHouse已确认写入。

1const { createClient } = require("@clickhouse/client");23classTickBatcher {4constructor() {5this.client = createClient({6url: process.env.CLICKHOUSE_URL || "http://localhost:8123",7username: process.env.CLICKHOUSE_USERNAME || "default",8password: process.env.CLICKHOUSE_PASSWORD || "",9database: process.env.CLICKHOUSE_DATABASE || "default",10compression: { request: true, response: true },11clickhouse_settings: { async_insert: 0 },12    });13this.batches = { trades: [], quotes: [] };14this.flushing = { trades: false, quotes: false };15this.failedRecords = 0;16this.droppedRecords = 0;17this.timer = setInterval(() => {18voidthis.flushBatch("trades");19voidthis.flushBatch("quotes");20    }, 2000);21  }2223handleMessage(data) {24const payload = JSON.parse(data.toString());25for (const { ev, ...fields } of payload) {26const table = ev === "T" ? "trades" : ev === "Q" ? "quotes" : null;27if (!table) continue;28if (this.batches[table].length >= 10000) {29this.droppedRecords++;30continue;31      }32this.batches[table].push(fields);33if (this.batches[table].length >= 1000) voidthis.flushBatch(table);34    }35  }3637asyncflushBatch(table) {38if (this.flushing[table]) return;39this.flushing[table] = true;40try {41while (this.batches[table].length) {42const values = this.batches[table].splice(0, 1000);43try {44awaitthis.client.insert({ table, values, format: "JSONEachRow" });45        } catch (error) {46this.failedRecords += values.length;47console.error(`Insert failed for ${table}:`, error.message);48        }49      }50    } finally {51this.flushing[table] = false;52    }53  }5455asyncclose() {56clearInterval(this.timer);57while (Object.values(this.flushing).some(Boolean)) {58awaitnewPromise((resolve) =>setTimeout(resolve, 10));59    }60awaitPromise.all([this.flushBatch("trades"), this.flushBatch("quotes")]);61awaitthis.client.close();62  }63}

紧凑类将每个待处理表缓冲区限制为10,000条记录,并计数丢弃或失败的记录;它不会重试失败的写入。对于无丢失管道,请添加持久缓冲、重放和去重。完整演示在1,000行或每两秒刷新一次,最多并发执行四次插入。

可视化实时市场数据

数据存储到ClickHouse后,构建可视化层就很简单了。主要挑战在于编写正确的SQL查询。让我们看看如何实现这一点。我们将重点介绍为两个关键可视化提供支持的查询。第一个是实时表格,它持续更新以显示特定股票的最新交易数据。

tick-table.png

为了构建此可视化,一个查询就足够了,数据可以使用ClickHouse强大的SQL查询语言和自定义函数进行格式化。

1with2{syms: Array(String)} as symbols,3toDate(now('America/New_York')) as curr_day,4trades_info as (5select6        sym,7        argMax(p, tuple(t, q)) as last_price,8        round(((last_price - (argMin(p, tuple(t, q)))) /nullIf(argMin(p, tuple(t, q)), 0)) *100, 2) as change_pct,9sum(s) as total_volume,10max(t) as latest_t11from12        trades13where14        toDate(fromUnixTimestamp64Milli(toInt64(t), 'America/New_York')) = curr_day15and sym in symbols16groupby17        sym18orderby19        sym asc20),21quotes_info as (22select23        sym,24        argMax(bp, tuple(t, q)) as bid,25        argMax(ap, tuple(t, q)) as ask,26max(t) as latest_t27from28        quotes29where30        toDate(fromUnixTimestamp64Milli(toInt64(t), 'America/New_York')) = curr_day31and sym in symbols32groupby33        sym34orderby35        sym asc36)37select38    t.sym as ticker,39    t.last_price aslast,40    q.bid as bid,41    q.ask as ask,42    t.change_pct as change,43    t.total_volume as volume,44    toUnixTimestamp64Milli(now64()) - toInt64(greatest(t.latest_t, q.latest_t)) as latency45from46    trades_info as t47leftjoin quotes_info as q on t.sym = q.sym;

让我们分解一下这个查询的作用。

首先,它定义了两个变量:symbols,其中包含要分析的股票代码列表,以及curr_day,它捕获纽约时区的当前日期。

然后查询检索交易数据,包括:

  • last_price:最新交易价格,使用argMax(p, tuple(t, q))以序列号打破毫秒时间戳并列
  • change_pct:相对于当前纽约日历日首次摄取交易的百分比变化,具有零分母保护。这不是官方开盘价或前收盘价。
  • total_volume:该日历日摄取的总整数股成交量

它还获取报价数据:

  • bid:最新买价,使用argMax(bp, tuple(t, q))
  • ask:最新卖价,使用argMax(ap, tuple(t, q))

最后,交易和报价结果合并。额外的latency列报告最新接收事件的年龄(以毫秒为单位),而不仅仅是数据库插入时间。以下行是说明性历史值。演示不筛选交易条件,不处理取消/更正,也不重建官方交易所OHLCV序列。

tickerlastbidaskchangevolume
NVDA151.2099151.2151.212.1765269276

我们要分析的第二个可视化是蜡烛图可视化,显示给定股票的价格演变和成交量。

让我们看看支持此可视化的SQL查询。

1select2toUnixTimestamp64Milli(toDateTime64(toStartOfInterval(fromUnixTimestamp64Milli(toInt64(t)), interval2minute), 3)) as x,3argMin(p, tuple(t, q)) as o,4max(p) as h,5min(p) as l,6argMax(p, tuple(t, q)) as c,7sum(s) as v8from trades9where x > toUnixTimestamp64Milli(now64() -interval1hour)10and sym = {sym: String}11groupby x orderby x asc;

此查询将最后一小时分组为两分钟桶,并计算开盘价、最高价、最低价、收盘价和成交量。开盘价和收盘价使用(t, q)解决时间戳相同的交易。桶过滤器排除了窗口开始时部分重叠的桶。

浏览器调用命名查询的Express端点;仅接受服务器定义的参数化SQL。ClickHouse凭据保留在服务器上,绝不放入NEXT_PUBLIC_变量中。为了可视化查询结果,我们使用click-ui组件进行表格显示,并使用Chart.js进行蜡烛图可视化。

扩展和实用技巧

在生产环境中处理高频市场数据不仅需要快速数据库。以下技巧和技术有助于确保系统在数据量增长时保持高性能和可靠性。

扩展摄取

处理许多代码的逐笔数据时,持续吞吐量可能轻松超过每秒数万条记录。

为处理此问题,需要注意不同方面:

  • 使用客户端批处理,插入大小根据系统的内存和延迟约束进行优化。
  • 使用压缩:压缩插入数据可减少网络发送的负载大小,最小化带宽使用并加速传输。
  • 监控ClickHouse中创建的分区数量,以防止过多合并。这篇博客讨论了异步插入,但分区创建部分可应用于同步摄取。您也可以使用高级仪表板监控数据分区数量。
  • Massive还提供了性能提示以处理高容量数据消费。

监控摄取延迟

如前所述,拥有最新数据对金融应用至关重要。因此监控它是合理的。

您可以轻松计算和跟踪事件时间戳(报价发生时间)与摄取时间戳(存储时间)之间的差异。此事件到插入延迟包括源传递、客户端批处理、网络时间和数据库摄取;它不是纯数据库延迟。下面的查询测量最近一小时每个代码最新交易的此延迟。有符号算术避免时间戳意外差异时的无符号下溢。

1SELECT2    sym,3count() AS trade_count,4    argMax(toInt64(inserted_at) - toInt64(t), tuple(t, q)) AS ingest_latency_ms5FROM trades6WHERE t >= toUnixTimestamp64Milli(now64() -INTERVAL1HOUR)7GROUPBY sym8ORDERBY trade_count DESC9LIMIT 100;

在仪表板中可视化此指标有助于您及早发现减速。

利用物化视图

物化视图当你希望在数据到达时进行预聚合时,这些很有用。这有助于优化依赖基于时间汇总的特定查询模式。一个典型例子是计算固定时间间隔(如1分钟窗口)内金融数据的OHLCV(开盘、最高、最低、收盘、成交量)指标。通过在摄取期间生成这些聚合,你可以快速提供结果,而无需每次重新计算。

首先创建一个目标表来存储1分钟OHLCV聚合。使用AggregatingMergeTree来合并聚合状态。将每个分组维度都包含在排序键中:在这里是symbol、tape和minute。开盘和收盘状态使用与原始查询相同的(t, q)排序。

1CREATE TABLE trades_1min_ohlcv2(3    `sym` LowCardinality(String),4    `z` Enum8('NYSE'=1, 'AMEX'=2, 'Nasdaq'=3),5    `minute_bucket_ms` UInt64,6    `open_price_state` AggregateFunction(argMin, Float64, Tuple(UInt64, UInt64)),7    `high_price_state` AggregateFunction(max, Float64),8    `low_price_state` AggregateFunction(min, Float64),9    `close_price_state` AggregateFunction(argMax, Float64, Tuple(UInt64, UInt64)),10    `volume_state` AggregateFunction(sum, UInt64),11    `trade_count_state` AggregateFunction(count)12)13ENGINE = AggregatingMergeTree14ORDERBY (sym, z, minute_bucket_ms);

下一步是创建物化视图。

1CREATE MATERIALIZED VIEW trades_1min_ohlcv_mv TO trades_1min_ohlcv2ASSELECT3    sym,4    z,5    intDiv(t, 60000) *60000AS minute_bucket_ms,6    argMinState(p, tuple(t, q)) AS open_price_state,7    maxState(p) AS high_price_state,8    minState(p) AS low_price_state,9    argMaxState(p, tuple(t, q)) AS close_price_state,10    sumState(s) AS volume_state,11    countState() AS trade_count_state12FROM trades13GROUPBY sym, z, minute_bucket_ms;

在开始摄取之前创建目标表和视图。视图处理新插入的数据块;它不会自动回填现有行。如果需要,使用协调的INSERT ... SELECT回填现有数据,注意不要重复计算重叠的实时数据。查询使用-Merge函数在后台合并完成之前组合状态。

要查看数据,请执行此查询。

1-- Query the table2SELECT3    sym,4    z,5    minute_bucket_ms,6    fromUnixTimestamp64Milli(toInt64(minute_bucket_ms)) as minute_timestamp,7    argMinMerge(open_price_state) AS open_price,8    maxMerge(high_price_state) AS high_price,9    minMerge(low_price_state) AS low_price,10    argMaxMerge(close_price_state) AS close_price,11    sumMerge(volume_state) AS volume,12    countMerge(trade_count_state) AS trade_count13FROM trades_1min_ohlcv14GROUPBY sym, z, minute_bucket_ms15ORDERBY sym, z, minute_bucket_ms;

运行完整示例

The example source 现在位于blog-examples/stock-data-demo下。使用Node.js 22或更高版本,以及现有的本地或云ClickHouse数据库。

1git clone https://github.com/ClickHouse/examples.git2cd examples/blog-examples/stock-data-demo3npm ci4cp env.example .env5# Set CLICKHOUSE_URL, CLICKHOUSE_USERNAME, CLICKHOUSE_PASSWORD,6# CLICKHOUSE_DATABASE, and MASSIVE_API_KEY in .env.7npm run setup8npm run build9npm start

打开http://localhost:34567/stocks/查看图表,或http://localhost:34567/admin进行摄取控制。默认情况下,服务订阅AAPL、MSFT和NVDA。更改MASSIVE_SYMBOLS以选择其他代码;仪表盘的观察列表不会改变上游订阅。服务器绑定到回环地址,其API没有身份验证,因此在公开之前添加访问控制。

如果你的密钥缺少必需的渠道,管理页面会报告身份验证失败。市场在周末和假日也可能很安静或关闭。要尝试没有市场访问权限的仪表盘,请使用单独的ClickHouse数据库,将API密钥留空,并先运行npm run seed,然后再运行npm start。这会插入带有过去三分钟时间戳的明确合成价格;它不会获取真实市场数据。当这些时间戳移动到图表窗口之外时,重新运行种子。

The README 涵盖了本地ClickHouse设置、开发模式、配置和验证。

结论

在这篇文章中,我们探讨了如何使用Massive获取市场数据和ClickHouse进行快速摄取和查询来构建实时逐笔数据应用。我们涵盖了如何流式处理和组织逐笔数据、管理摄取性能,以及构建高效的查询和可视化。

在这个GitHubrepository中,你会找到一个使用React作为可视化层的可运行示例。用它来探索摄取和查询模式,然后为生产应用添加持久交付、重放、认证和针对工作负载的调优。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存