|
|
|
|
@ -40,7 +40,7 @@ type KlineSeries struct {
|
|
|
|
|
Interval types.Interval |
|
|
|
|
IntervalAdder types.IntervalAdder |
|
|
|
|
lastTs int64 |
|
|
|
|
klines []*types.Kline |
|
|
|
|
ringSeries *types.RingSeries[*types.Kline] |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func NewKlineSeries(exchange pb.ExchangeType, instId string, interval types.Interval) *KlineSeries { |
|
|
|
|
@ -53,7 +53,7 @@ func NewKlineSeries(exchange pb.ExchangeType, instId string, interval types.Inte
|
|
|
|
|
InstId: instId, |
|
|
|
|
Interval: interval, |
|
|
|
|
IntervalAdder: intervalAdder, |
|
|
|
|
klines: make([]*types.Kline, 0, MaxSeriesKlines/10), |
|
|
|
|
ringSeries: types.NewRingSeries[*types.Kline](MaxSeriesKlines, MaxSeriesKlines/10), |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -66,8 +66,117 @@ func (s *KlineSeries) MustGet(offset int16) (k types.Kline) {
|
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// GetE [0]当前k线
|
|
|
|
|
// Get [0]当前k线
|
|
|
|
|
func (s *KlineSeries) Get(offset int16) (k types.Kline, err error) { |
|
|
|
|
s.mu.RLock() |
|
|
|
|
defer s.mu.RUnlock() |
|
|
|
|
v, ok := s.ringSeries.Get(int(offset)) |
|
|
|
|
if !ok { |
|
|
|
|
err = fmt.Errorf("get kline series offset out of range: offset=%d, length=%d", offset, s.ringSeries.Length()) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
k = *v |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) MustSeries(offset, count int16) (klines series.Klines) { |
|
|
|
|
klines, err := s.Series(offset, count) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Series 时间降序序列[count...offset]
|
|
|
|
|
// offset: 从序列尾部开始偏移量
|
|
|
|
|
// count: 从offset位置开始向序列头部k线条数
|
|
|
|
|
func (s *KlineSeries) Series(offset, count int16) (klines series.Klines, err error) { |
|
|
|
|
s.mu.RLock() |
|
|
|
|
defer s.mu.RUnlock() |
|
|
|
|
series, ok := s.ringSeries.Series(int(offset), int(count)) |
|
|
|
|
if !ok { |
|
|
|
|
err = fmt.Errorf("get kline series offset out of range: offset=%d, count=%d", offset, count) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
for _, si := range series { |
|
|
|
|
klines = append(klines, *si) |
|
|
|
|
} |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) Length() int { |
|
|
|
|
s.mu.RLock() |
|
|
|
|
defer s.mu.RUnlock() |
|
|
|
|
return s.ringSeries.Length() |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) LastTs() int64 { |
|
|
|
|
return s.lastTs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// 检查k线序列完整
|
|
|
|
|
func (s *KlineSeries) Update(kline *types.Kline) (lastTs int64, serial bool) { |
|
|
|
|
s.mu.Lock() |
|
|
|
|
defer s.mu.Unlock() |
|
|
|
|
|
|
|
|
|
serial = true |
|
|
|
|
lastTs = s.lastTs |
|
|
|
|
if kline.Ts <= s.lastTs { |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// 检查k线是否连续
|
|
|
|
|
if s.lastTs != 0 { |
|
|
|
|
expectTs := s.IntervalAdder(s.lastTs, 1) |
|
|
|
|
if kline.Ts != expectTs { |
|
|
|
|
serial = false |
|
|
|
|
miss := (kline.Ts-s.lastTs)/s.IntervalAdder(0, 1) - 1 |
|
|
|
|
zlog.Warningf("k线不连续: instId=%s(%s), interval=%s, miss=%d, lastTs=%d, got=%d, expected=%d", s.InstId, s.Exchange, s.Interval, miss, s.lastTs, kline.Ts, expectTs) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
s.ringSeries.Push(kline) |
|
|
|
|
s.lastTs = kline.Ts |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Deprecated: use KlineSeries
|
|
|
|
|
type KlineSeries0 struct { |
|
|
|
|
mu sync.RWMutex |
|
|
|
|
Exchange pb.ExchangeType |
|
|
|
|
InstId string |
|
|
|
|
Interval types.Interval |
|
|
|
|
IntervalAdder types.IntervalAdder |
|
|
|
|
lastTs int64 |
|
|
|
|
klines []*types.Kline |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func NewKlineSeries0(exchange pb.ExchangeType, instId string, interval types.Interval) *KlineSeries0 { |
|
|
|
|
intervalAdder, ok := types.SupportedIntervals[interval] |
|
|
|
|
if !ok { |
|
|
|
|
panic(fmt.Errorf("unsupport interval: %s", interval)) |
|
|
|
|
} |
|
|
|
|
return &KlineSeries0{ |
|
|
|
|
Exchange: exchange, |
|
|
|
|
InstId: instId, |
|
|
|
|
Interval: interval, |
|
|
|
|
IntervalAdder: intervalAdder, |
|
|
|
|
klines: make([]*types.Kline, 0, MaxSeriesKlines/10), |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Get [0]当前k线
|
|
|
|
|
func (s *KlineSeries0) MustGet(offset int16) (k types.Kline) { |
|
|
|
|
k, err := s.Get(offset) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// GetE [0]当前k线
|
|
|
|
|
func (s *KlineSeries0) Get(offset int16) (k types.Kline, err error) { |
|
|
|
|
if ok := offset >= 0 && offset < MaxSeriesKlines; !ok { |
|
|
|
|
err = fmt.Errorf("get kline series offset out of range: offset=%d", offset) |
|
|
|
|
return |
|
|
|
|
@ -84,7 +193,7 @@ func (s *KlineSeries) Get(offset int16) (k types.Kline, err error) {
|
|
|
|
|
return *(s.klines[index]), nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) MustSeries(offset, count int16) (klines series.Klines) { |
|
|
|
|
func (s *KlineSeries0) MustSeries(offset, count int16) (klines series.Klines) { |
|
|
|
|
klines, err := s.Series(offset, count) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
@ -95,7 +204,7 @@ func (s *KlineSeries) MustSeries(offset, count int16) (klines series.Klines) {
|
|
|
|
|
// Series 时间降序序列[count...offset]
|
|
|
|
|
// offset: 从序列尾部开始偏移量
|
|
|
|
|
// count: 从offset位置开始向序列头部k线条数
|
|
|
|
|
func (s *KlineSeries) Series(offset, count int16) (klines series.Klines, err error) { |
|
|
|
|
func (s *KlineSeries0) Series(offset, count int16) (klines series.Klines, err error) { |
|
|
|
|
if ok := offset >= 0 && offset < MaxSeriesKlines; !ok { |
|
|
|
|
err = fmt.Errorf("get kline series offset out of range: offset=%d, count=%d", offset, count) |
|
|
|
|
return |
|
|
|
|
@ -124,18 +233,18 @@ func (s *KlineSeries) Series(offset, count int16) (klines series.Klines, err err
|
|
|
|
|
return klines, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) Length() int { |
|
|
|
|
func (s *KlineSeries0) Length() int { |
|
|
|
|
s.mu.RLock() |
|
|
|
|
defer s.mu.RUnlock() |
|
|
|
|
return len(s.klines) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *KlineSeries) LastTs() int64 { |
|
|
|
|
func (s *KlineSeries0) LastTs() int64 { |
|
|
|
|
return s.lastTs |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// 检查k线序列完整
|
|
|
|
|
func (s *KlineSeries) Update(kline *types.Kline) (lastTs int64, serial bool) { |
|
|
|
|
func (s *KlineSeries0) Update(kline *types.Kline) (lastTs int64, serial bool) { |
|
|
|
|
s.mu.Lock() |
|
|
|
|
defer s.mu.Unlock() |
|
|
|
|
|
|
|
|
|
|