基于Node.js构建加密货币价格实时告警系统的技术方案咨询
基于Node.js实现加密货币实时价格告警的优化方案
核心思路:事件驱动替代轮询
你之前用setInterval轮询的问题在于,不管价格有没有变化都要遍历所有规则,既浪费资源又无法做到真正实时。换成价格更新事件触发规则匹配的模式,能从根本上解决实时性问题。
1. 订阅实时价格流
直接对接交易所的WebSocket API,实时接收币种价格变动。每次收到新价格时才触发规则校验,完全贴合价格变化的时间点,无需等待轮询周期。
2. 优化规则存储与匹配
针对1000条规则,通过结构化存储和高效查找避免全量遍历:
- 按币种+触发类型分组:比如将BTC的“价格高于阈值”规则存入一个数组,“价格低于阈值”规则存入另一个数组。
- 对每组规则的阈值排序:“高于阈值”数组按从小到大排序,“低于阈值”数组按从大到小排序。
- 用二分查找快速定位触发规则:收到新价格后,在对应数组中快速筛选出符合条件的规则,无需遍历全部。
示例代码片段:
// 规则存储结构(按币种和触发类型分组,阈值已排序) const alertRules = { "BTC": { above: [ { userId: 1, threshold: 50000, notified: false }, { userId: 2, threshold: 51000, notified: false } ], below: [ { userId: 3, threshold: 48000, notified: false }, { userId: 4, threshold: 47000, notified: false } ] } }; // 处理价格更新 function handlePriceUpdate(symbol, newPrice, lastPrice) { const rules = alertRules[symbol]; if (!rules) return; // 处理价格上涨触发的"高于阈值"规则 if (newPrice > lastPrice) { const triggerIndex = findFirstGreater(rules.above, newPrice); rules.above.slice(0, triggerIndex).forEach(rule => { if (!rule.notified) { addToNotificationQueue(rule.userId, `${symbol}价格突破${rule.threshold}`); rule.notified = true; // 标记已通知,避免重复触发 } }); } // 处理价格下跌触发的"低于阈值"规则 if (newPrice < lastPrice) { const triggerIndex = findFirstSmaller(rules.below, newPrice); rules.below.slice(0, triggerIndex).forEach(rule => { if (!rule.notified) { addToNotificationQueue(rule.userId, `${symbol}价格跌破${rule.threshold}`); rule.notified = true; } }); } } // 二分查找:找到第一个阈值大于目标价格的索引 function findFirstGreater(arr, target) { let low = 0, high = arr.length; while (low < high) { const mid = Math.floor((low + high) / 2); if (arr[mid].threshold > target) { high = mid; } else { low = mid + 1; } } return low; } // 二分查找:找到第一个阈值小于目标价格的索引 function findFirstSmaller(arr, target) { let low = 0, high = arr.length; while (low < high) { const mid = Math.floor((low + high) / 2); if (arr[mid].threshold < target) { high = mid; } else { low = mid + 1; } } return low; }
3. 异步队列处理通知
触发告警后,不要直接同步发送邮件/短信,把通知任务放到异步队列(比如bullmq),由单独的worker进程处理。这样不会阻塞价格处理的主线程,保证实时性。
示例队列代码:
const { Queue, Worker } = require('bullmq'); const notificationQueue = new Queue('alerts'); // 将通知任务加入队列 function addToNotificationQueue(userId, message) { notificationQueue.add('send-alert', { userId, message }); } // Worker进程处理通知任务 const worker = new Worker('alerts', async (job) => { const { userId, message } = job.data; // 这里实现具体的通知逻辑:调用邮件API、推送短信等 await sendNotification(userId, message); });
4. 避免单线程阻塞
- 价格处理逻辑尽量轻量化,不要在WebSocket回调中做耗时操作。
- 如果后续规则数量增长到数万级,可将规则匹配逻辑放到
worker_threads线程中执行,避免阻塞主线程的WebSocket连接。
额外注意事项
- 去重触发:记录每个规则的触发状态,或设置冷却时间,避免价格在阈值附近波动时重复发送告警。
- 规则动态更新:用户新增/修改/删除规则时,实时更新存储数组并重新排序,确保匹配准确性。
内容的提问来源于stack exchange,提问作者Tirth Sidhpura
相关产品推荐
相关产品推荐

