diff --git a/internal/trading/kline_series.go b/internal/trading/kline_series.go index 035553c..881595e 100644 --- a/internal/trading/kline_series.go +++ b/internal/trading/kline_series.go @@ -1,39 +1,53 @@ package trading import ( + "fmt" "sig-pub/api/pb" "sig-pub/pkg/types" "sig-pub/pkg/types/series" + "sig-pub/pkg/zlog" + "sync" +) + +const ( + MaxSeriesKlines = 1000 ) type KlineSeries struct { - Exchange pb.ExchangeType - InstId string - Interval types.Interval - Ts int64 - klines []*types.Kline - klineStore *KlineStore + sync.RWMutex + Exchange pb.ExchangeType + InstId string + Interval types.Interval + IntervalAdder types.IntervalAdder + lastTs int64 + klines []*types.Kline } -func NewKlineSeries(ts int64, interval types.Interval, klineStore *KlineStore) *KlineSeries { +func NewKlineSeries(exchange pb.ExchangeType, instId string, interval types.Interval) *KlineSeries { + intervalAdder, ok := types.SupportedIntervals[interval] + if !ok { + panic(fmt.Errorf("unsupport interval: %s", interval)) + } return &KlineSeries{ - Ts: ts, - Interval: interval, - klineStore: klineStore, + Exchange: exchange, + InstId: instId, + Interval: interval, + IntervalAdder: intervalAdder, + klines: make([]*types.Kline, 0, MaxSeriesKlines/10), } } // Get // [0]当前k线 -func (a KlineSeries) Get(start int16) types.Kline { - index := len(a.klines) - 1 - int(start) - if index >= 0 && index < len(a.klines)-1 { - return *(a.klines[index]) +func (s *KlineSeries) Get(start int16) types.Kline { + index := len(s.klines) - 1 - int(start) + if index >= 0 && index < len(s.klines)-1 { + return *(s.klines[index]) } // todo query store - ts := a.Interval.MustAddMul(a.Ts, int64(-start)) - for _, k := range a.klines { + ts := s.Interval.MustAddMul(s.lastTs, int64(-start)) + for _, k := range s.klines { if k.Ts == ts { return *k } @@ -43,9 +57,9 @@ func (a KlineSeries) Get(start int16) types.Kline { } // Series [start...end] -func (a KlineSeries) Series(start, end int16) (klines series.Klines) { - endTs := a.Interval.MustAddMul(a.Ts, int64(-start)) - startTs := a.Interval.MustAddMul(a.Ts, int64(-end)) +func (s *KlineSeries) Series(start, end int16) (klines series.Klines) { + endTs := s.Interval.MustAddMul(s.lastTs, int64(-start)) + startTs := s.Interval.MustAddMul(s.lastTs, int64(-end)) _ = endTs _ = startTs // return a.klineStore.GetRange(startTs, endTs) @@ -54,9 +68,38 @@ func (a KlineSeries) Series(start, end int16) (klines series.Klines) { } // 检查k线序列完整 -func (s *KlineSeries) Update(kline *types.Kline) []types.Kline { - s.klines = append(s.klines, kline) - // todo copy(s.klines, s.klines[0:1]) set index=20, ts=kline.ts - s.Ts = kline.Ts - return nil +func (s *KlineSeries) Update(kline *types.Kline) (lastTs int64, serial bool) { + s.Lock() + defer s.Unlock() + + serial = true + lastTs = s.lastTs + if kline.Ts <= s.lastTs { + // 已存在的k线,直接忽略 + return + } + // 检查k线是否连续 + if len(s.klines) > 0 { + expectTs := s.Interval.MustAddMul(s.lastTs, 1) + if kline.Ts != expectTs { + + serial = false + zlog.Warningf("k线不连续: last.Ts=%d, expected=%d, got=%d", s.lastTs, expectTs, kline.Ts) + return + } + } + + if cap(s.klines) < MaxSeriesKlines { + s.klines = append(s.klines, kline) + } else { + if len(s.klines) < MaxSeriesKlines { + s.klines = append(s.klines, kline) + } else { + // 循环复用切片空间,避免扩容 + copy(s.klines, s.klines[1:]) + s.klines[MaxSeriesKlines-1] = kline + } + } + s.lastTs = kline.Ts + return } diff --git a/internal/trading/kline_store.go b/internal/trading/kline_store.go index beff5b4..4df07c5 100644 --- a/internal/trading/kline_store.go +++ b/internal/trading/kline_store.go @@ -9,24 +9,60 @@ import ( type KlineStore struct { vmdb vmts.VictoriaMetricsTSDB - store [3]*collect.ConcurrentMap[string, *collect.ConcurrentMap[types.Interval, *KlineSeries]] // K线列表: []exchange + store [3]*collect.ConcurrentMap[string, *KlineStoreInstance] // K线列表: []exchange } func NewKlineSeriesStore() (kss *KlineStore) { kss = &KlineStore{} - kss.store[pb.ExchangeType_OKX] = collect.NewConcurrentMap[string, *collect.ConcurrentMap[types.Interval, *KlineSeries]](64, func(s string) string { return s }) + kss.store[pb.ExchangeType_OKX] = collect.NewConcurrentMap[string, *KlineStoreInstance](64, func(s string) string { return s }) // kss.klines[pb.ExchangeType_BINANCE] = return } func (s *KlineStore) Update(exchange pb.ExchangeType, instId string, kline *types.Kline) (k types.Kline) { - // klineSeries Update - intervals := s.store[exchange].ComputeIfAbsent(instId, func(k string) *collect.ConcurrentMap[types.Interval, *KlineSeries] { - return collect.NewConcurrentMap[types.Interval, *KlineSeries](8, func(t types.Interval) string { return string(t) }) + storeInst := s.store[exchange].ComputeIfAbsent(instId, func(k string) *KlineStoreInstance { + return NewKlineStoreInstance(exchange, k) }) - klineSeries := intervals.ComputeIfAbsent(kline.Interval, func(interval types.Interval) *KlineSeries { - return NewKlineSeries(0, interval, s) - }) - klineSeries.Update(kline) + storeInst.Update(kline) + return +} + +// KlineStoreInstance 单个交易产品所有周期k线 +type KlineStoreInstance struct { + Exchange pb.ExchangeType + InstId string + IntervalKlines *collect.SyncMap[types.Interval, *KlineSeries] +} + +func NewKlineStoreInstance(exchange pb.ExchangeType, instId string) *KlineStoreInstance { + si := &KlineStoreInstance{ + Exchange: exchange, + InstId: instId, + IntervalKlines: collect.NewSyncMap[types.Interval, *KlineSeries](), + } + for interval := range types.SupportedIntervals { + si.IntervalKlines.Store(interval, NewKlineSeries(exchange, instId, interval)) + } + return si +} + +// Update 更新k线 +// kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series +func (si *KlineStoreInstance) Update(kline *types.Kline) (ok bool) { + ks, ok := si.IntervalKlines.Load(kline.Interval) + if !ok { + return + } + + if lastTs, serial := ks.Update(kline); !serial { + // k线不完整 + // 拉取k线 + start := lastTs + end := kline.Ts + _, _ = start, end + } + + ok = true + // emit kline event return } diff --git a/internal/trading/trading_service.go b/internal/trading/trading_service.go index d0689bd..efb5d4c 100644 --- a/internal/trading/trading_service.go +++ b/internal/trading/trading_service.go @@ -111,6 +111,7 @@ func (svr *TradingService) subscribeKlines(reconnect bool) { kline := new(types.Kline) kline.ParsePBKline(msg.Kline.Exchange, k) zlog.Debugf("recv: streamId=%d, %v, %s, %#v", msg.Kline.StreamId, msg.Kline.Exchange, msg.Kline.InstId, kline) + // kline klineStore -> klineSeries -> strategy -> indicator -> klineSeries.Series svr.klineStore.Update(msg.Kline.Exchange, msg.Kline.InstId, kline) } }