寻求固定容量、自动移除最早元素的集合用于均线指标计算
问题解答
Channels是否适用?
不适用。Channels
推荐实现方案
最直接的方式是基于Queue<T>自定义一个固定容量的滑动窗口集合,或者使用现成的环形缓冲区实现。以下是两种实用方案:
方案1:自定义固定窗口集合
实现一个简单的FixedSizeWindow<T>类,封装Queue
public class FixedSizeWindow<T> { private readonly Queue<T> _queue; public int Capacity { get; } public int Count => _queue.Count; public FixedSizeWindow(int capacity) { Capacity = capacity > 0 ? capacity : throw new ArgumentOutOfRangeException(nameof(capacity)); _queue = new Queue<T>(capacity); } public void Add(T item) { _queue.Enqueue(item); // 容量满时移除最早元素 while (_queue.Count > Capacity) { _queue.Dequeue(); } } // 提供获取所有元素的方法,方便计算指标 public IEnumerable<T> GetAllItems() => _queue.ToArray(); }
方案2:使用Rx.NET的Buffer操作符(适合响应式场景)
如果你的代码本来就基于响应式编程,Rx.NET的Buffer操作符可以直接实现滑动窗口逻辑,自动维护指定数量的最新元素:
// 假设你已经将Trade数据流转换为IObservable<Trade> tradeObservable.Buffer(loopback, 1) // 窗口大小loopback,每次滑动1个元素 .Subscribe(window => { if (window.Count == loopback) { // 计算SMA/EMA var sma = window.Average(t => t.Price); // 发布信号 } }, stoppingToken);
整合到你的现有代码
使用方案1的自定义集合,修改你的代码如下:
// 初始化固定窗口集合,loopback是你的回溯周期大小 var tradeWindow = new FixedSizeWindow<Trade>(loopback); var subResult = await _socketClient.SpotStreams.SubscribeToTradeUpdatesAsync(symbols, data => { var trade = new Trade(data.Timestamp, data.Data.Symbol, data.Data.Price, data.Data.Quantity); // 添加到窗口集合,自动维护容量 tradeWindow.Add(trade); // 当窗口元素数量达到回溯周期时计算指标 if (tradeWindow.Count == tradeWindow.Capacity) { var allTrades = tradeWindow.GetAllItems(); // 计算SMA:平均价格 var sma = allTrades.Average(t => t.Price); // 计算EMA:指数加权逻辑示例 var ema = CalculateEma(allTrades.Select(t => t.Price)); // 生成信号并处理 // signal = ...; } publisher.Publish(trade); }, stoppingToken); // 示例EMA计算方法 private decimal CalculateEma(IEnumerable<decimal> prices) { var priceList = prices.ToList(); var alpha = 2m / (priceList.Count + 1); var ema = priceList[0]; for (int i = 1; i < priceList.Count; i++) { ema = alpha * priceList[i] + (1 - alpha) * ema; } return ema; }
内容的提问来源于stack exchange,提问作者nop
相关产品推荐
相关产品推荐

