Browse Source

kline inst series update

main
strange 11 months ago
parent
commit
f580fe1da0
  1. 79
      internal/trading/kline_series.go
  2. 54
      internal/trading/kline_store.go
  3. 1
      internal/trading/trading_service.go

79
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 {
sync.RWMutex
Exchange pb.ExchangeType
InstId string
Interval types.Interval
Ts int64
IntervalAdder types.IntervalAdder
lastTs int64
klines []*types.Kline
klineStore *KlineStore
}
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,
Exchange: exchange,
InstId: instId,
Interval: interval,
klineStore: klineStore,
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 {
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)
// todo copy(s.klines, s.klines[0:1]) set index=20, ts=kline.ts
s.Ts = kline.Ts
return nil
} else {
// 循环复用切片空间,避免扩容
copy(s.klines, s.klines[1:])
s.klines[MaxSeriesKlines-1] = kline
}
}
s.lastTs = kline.Ts
return
}

54
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<instId, interval, klines>
store [3]*collect.ConcurrentMap[string, *KlineStoreInstance] // K线列表: []exchange<instId, interval, klines>
}
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
}

1
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)
}
}

Loading…
Cancel
Save