Browse Source

okx kline fetch

main
strange 1 year ago
parent
commit
5d802be0fd
  1. 10
      cmd/test/test.go
  2. 4
      go.mod
  3. 4
      go.sum
  4. 34
      internal/exchange/exchange_grpc_server.go
  5. 19
      internal/exchange/okx/channel_kline.go
  6. 25
      internal/exchange/okx/okx.go
  7. 92
      internal/exchange/okx/okx_fetch.go
  8. 44
      internal/exchange/okx/okx_fetch_test.go
  9. 12
      pkg/storage/tsdb/victoria_metrics/metric.go
  10. 68
      pkg/storage/tsdb/victoria_metrics/vm.go
  11. 2
      pkg/storage/tsdb/victoria_metrics/vm_test.go
  12. 8
      pkg/types/instance.go
  13. 89
      pkg/types/interval.go

10
cmd/test/test.go

@ -1,8 +1,6 @@
package main package main
import ( import (
"sig-pub/api/pb"
"sig-pub/pkg/types"
"sig-pub/pkg/zlog" "sig-pub/pkg/zlog"
"github.com/VictoriaMetrics/metrics" "github.com/VictoriaMetrics/metrics"
@ -24,10 +22,10 @@ type BTC struct {
func testDecimal() { func testDecimal() {
// zlog.Init() // zlog.Init()
var interval = "1m" // var interval = "1m"
interval0 := types.Interval(interval) // interval0 := types.Interval(interval)
sec, _ := interval0.Seconds() // sec, _ := interval0.Seconds()
zlog.Infof("hello...: %s, %s, %d", pb.Exchange_OKX.String(), pb.Exchange_BINANCE.String(), sec) // zlog.Infof("hello...: %s, %s, %d", pb.Exchange_OKX.String(), pb.Exchange_BINANCE.String())
// data := `{"price":"1.23456"}` // data := `{"price":"1.23456"}`
// btc := new(BTC) // btc := new(BTC)

4
go.mod

@ -9,6 +9,7 @@ require (
github.com/dsnet/golib/unitconv v1.0.2 github.com/dsnet/golib/unitconv v1.0.2
github.com/fanjindong/go-cache v0.0.6 github.com/fanjindong/go-cache v0.0.6
github.com/gin-gonic/gin v1.10.0 github.com/gin-gonic/gin v1.10.0
github.com/go-resty/resty/v2 v2.16.5
github.com/gorilla/websocket v1.5.3 github.com/gorilla/websocket v1.5.3
github.com/govalues/decimal v0.1.36 github.com/govalues/decimal v0.1.36
github.com/influxdata/influxdb-client-go/v2 v2.14.0 github.com/influxdata/influxdb-client-go/v2 v2.14.0
@ -21,6 +22,8 @@ require (
go.etcd.io/etcd/client/v3 v3.6.1 go.etcd.io/etcd/client/v3 v3.6.1
go.uber.org/zap v1.27.0 go.uber.org/zap v1.27.0
golang.org/x/net v0.38.0 golang.org/x/net v0.38.0
golang.org/x/sync v0.12.0
golang.org/x/time v0.8.0
google.golang.org/grpc v1.71.1 google.golang.org/grpc v1.71.1
google.golang.org/protobuf v1.36.6 google.golang.org/protobuf v1.36.6
gopkg.in/natefinch/lumberjack.v2 v2.2.1 gopkg.in/natefinch/lumberjack.v2 v2.2.1
@ -79,7 +82,6 @@ require (
go.uber.org/multierr v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect
golang.org/x/arch v0.15.0 // indirect golang.org/x/arch v0.15.0 // indirect
golang.org/x/crypto v0.36.0 // indirect golang.org/x/crypto v0.36.0 // indirect
golang.org/x/sync v0.12.0 // indirect
golang.org/x/sys v0.31.0 // indirect golang.org/x/sys v0.31.0 // indirect
golang.org/x/text v0.23.0 // indirect golang.org/x/text v0.23.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20250303144028-a0af3efb3deb // indirect google.golang.org/genproto/googleapis/api v0.0.0-20250303144028-a0af3efb3deb // indirect

4
go.sum

@ -57,6 +57,8 @@ github.com/go-playground/universal-translator v0.18.1 h1:Bcnm0ZwsGyWbCzImXv+pAJn
github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY= github.com/go-playground/universal-translator v0.18.1/go.mod h1:xekY+UJKNuX9WP91TpwSH2VMlDf28Uj24BCp08ZFTUY=
github.com/go-playground/validator/v10 v10.26.0 h1:SP05Nqhjcvz81uJaRfEV0YBSSSGMc/iMaVtFbr3Sw2k= github.com/go-playground/validator/v10 v10.26.0 h1:SP05Nqhjcvz81uJaRfEV0YBSSSGMc/iMaVtFbr3Sw2k=
github.com/go-playground/validator/v10 v10.26.0/go.mod h1:I5QpIEbmr8On7W0TktmJAumgzX4CA1XNl4ZmDuVHKKo= github.com/go-playground/validator/v10 v10.26.0/go.mod h1:I5QpIEbmr8On7W0TktmJAumgzX4CA1XNl4ZmDuVHKKo=
github.com/go-resty/resty/v2 v2.16.5 h1:hBKqmWrr7uRc3euHVqmh1HTHcKn99Smr7o5spptdhTM=
github.com/go-resty/resty/v2 v2.16.5/go.mod h1:hkJtXbA2iKHzJheXYvQ8snQES5ZLGKMwQ07xAwp/fiA=
github.com/go-sql-driver/mysql v1.7.0 h1:ueSltNNllEqE3qcWBTD0iQd3IpL/6U+mJxLkazJ7YPc= github.com/go-sql-driver/mysql v1.7.0 h1:ueSltNNllEqE3qcWBTD0iQd3IpL/6U+mJxLkazJ7YPc=
github.com/go-sql-driver/mysql v1.7.0/go.mod h1:OXbVy3sEdcQ2Doequ6Z5BW6fXNQTmx+9S1MCJN5yJMI= github.com/go-sql-driver/mysql v1.7.0/go.mod h1:OXbVy3sEdcQ2Doequ6Z5BW6fXNQTmx+9S1MCJN5yJMI=
github.com/go-viper/mapstructure/v2 v2.2.1 h1:ZAaOCxANMuZx5RCeg0mBdEZk7DZasvvZIxtHqx8aGss= github.com/go-viper/mapstructure/v2 v2.2.1 h1:ZAaOCxANMuZx5RCeg0mBdEZk7DZasvvZIxtHqx8aGss=
@ -225,6 +227,8 @@ golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY= golang.org/x/text v0.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY=
golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4= golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4=
golang.org/x/time v0.8.0 h1:9i3RxcPv3PZnitoVGMPDKZSq1xW1gK1Xy3ArNOGZfEg=
golang.org/x/time v0.8.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE=

34
internal/exchange/exchange_grpc_server.go

@ -97,6 +97,11 @@ func (svc *ExchangeGrpcServer) subscribeExchanges() {
// consumerKline 消费交易所k线数据 // consumerKline 消费交易所k线数据
func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types.ChannelKline) { func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types.ChannelKline) {
exchangeType, ok := types.ExchangePBParse(exchange.Type)
if !ok {
panic(fmt.Errorf("unknown exchange type: %v", exchange.Type))
}
for { for {
channelK, ok := <-c channelK, ok := <-c
if !ok { if !ok {
@ -115,18 +120,6 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
continue continue
} }
// tsdb storage
typeInst := types.TradeInstance{
InstId: exInst.InstId, // channelK.InstId
TickSz: 0,
MinSz: 0,
}
// todo 异步处理
err := svc.exchangeDataService.SaveKlines(typeInst, channelK.Klines)
if err != nil {
zlog.Errorf("kline save to tsdb error: ", err)
}
// publish to subscribers // publish to subscribers
pubMsgMap := make(map[string]*pb.StreamKline) pubMsgMap := make(map[string]*pb.StreamKline)
@ -138,11 +131,13 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
} }
exchangeName := exchange.String() exchangeName := exchange.String()
var confirmKlines []*types.Kline
for _, kline := range channelK.Klines { for _, kline := range channelK.Klines {
// zlog.Infof("recv kline: %#v", kline) // zlog.Infof("recv kline: %#v", kline)
confirm := 0 confirm := 0
if kline.Confirm { if kline.Confirm {
confirm = 1 confirm = 1
confirmKlines = append(confirmKlines, kline)
} }
pubKey := fmt.Sprintf("/kline/%s/%s/%s/%d", exchangeName, exInst.InstId, kline.Interval, confirm) pubKey := fmt.Sprintf("/kline/%s/%s/%s/%d", exchangeName, exInst.InstId, kline.Interval, confirm)
// todo 优化没有订阅者就跳过 // todo 优化没有订阅者就跳过
@ -158,6 +153,21 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
msg.Klines = append(msg.Klines, pbk) msg.Klines = append(msg.Klines, pbk)
} }
// tsdb storage
if len(confirmKlines) > 0 {
typeInst := types.TradeInstance{
InstId: exInst.InstId, // channelK.InstId
PriceSz: 0,
QuantitySz: 0,
Exchange: exchangeType,
}
// todo 异步处理
err := svc.exchangeDataService.SaveKlines(typeInst, confirmKlines)
if err != nil {
zlog.Errorf("kline save to tsdb error: ", err)
}
}
for pubKey, msg := range pubMsgMap { for pubKey, msg := range pubMsgMap {
if len(msg.Klines) == 0 { if len(msg.Klines) == 0 {
continue continue

19
internal/exchange/okx/channel_kline.go

@ -20,13 +20,14 @@ var (
// UTC时间开盘价k线:[6Hutc/12Hutc/1Dutc/2Dutc/3Dutc/1Wutc/1Mutc/3Mutc] // UTC时间开盘价k线:[6Hutc/12Hutc/1Dutc/2Dutc/3Dutc/1Wutc/1Mutc/3Mutc]
// "candle1s" candle3Mutc candle1Mutc candle1Wutc candle1Dutc candle2Dutc candle3Dutc candle5Dutc candle12Hutc candle6Hutc // "candle1s" candle3Mutc candle1Mutc candle1Wutc candle1Dutc candle2Dutc candle3Dutc candle5Dutc candle12Hutc candle6Hutc
// candles = []string{"candle3M", "candle1M", "candle1W", "candle1D", "candle2D", "candle3D", "candle5D", "candle12H", "candle6H", "candle4H", "candle2H", "candle1H", "candle30m", "candle15m", "candle5m", "candle3m", "candle1m"} // candles = []string{"candle3M", "candle1M", "candle1W", "candle1D", "candle2D", "candle3D", "candle5D", "candle12H", "candle6H", "candle4H", "candle2H", "candle1H", "candle30m", "candle15m", "candle5m", "candle3m", "candle1m"}
subCandles = []string{ //
// subCandles = []string{
// "candle3M", "candle1M", // "candle3M", "candle1M",
"candle1W", "candle1D", "candle2D", "candle3D", "candle5D", // "candle1W", "candle1D", "candle2D", "candle3D", "candle5D",
"candle12H", "candle6H", "candle4H", "candle2H", "candle1H", // "candle12H", "candle6H", "candle4H", "candle2H", "candle1H",
"candle30m", "candle15m", "candle5m", "candle3m", "candle1m", // "candle30m", "candle15m", "candle5m", "candle3m", "candle1m",
"candle1s", // "candle1s",
} // }
) )
// ChannelCandle k线订阅频道 // ChannelCandle k线订阅频道
@ -39,11 +40,15 @@ func NewChannelCandle(
traceId string, traceId string,
httpProxy string, httpProxy string,
) *ChannelCandle { ) *ChannelCandle {
var candles []string
for _, interval := range subscribeCandles {
candles = append(candles, "candle"+interval)
}
cfg := wsChannelConfig[*CandleData, *types.ChannelKline]{ cfg := wsChannelConfig[*CandleData, *types.ChannelKline]{
channelId: traceId, channelId: traceId,
httpProxy: httpProxy, httpProxy: httpProxy,
wsUrl: "/ws/v5/business", wsUrl: "/ws/v5/business",
subscribeChannels: subCandles, subscribeChannels: candles,
dataInstanceFunc: func() *CandleData { dataInstanceFunc: func() *CandleData {
var d CandleData var d CandleData
return &d return &d

25
internal/exchange/okx/okx.go

@ -5,6 +5,31 @@ import (
"sig-pub/pkg/types" "sig-pub/pkg/types"
) )
var subscribeCandles = map[types.Interval]string{
// "3M", "1M", "1W", "1D", "2D", "3D", "5D",
// "12H", "6H", "4H", "2H", "1H",
// "30m", "15m", "5m", "3m", "1m",
// "1s",
types.Interval1s: "1s",
types.Interval1m: "1m",
types.Interval3m: "3m",
types.Interval5m: "5m",
types.Interval15m: "15m",
types.Interval30m: "30m",
types.Interval1h: "1H",
types.Interval2h: "2H",
types.Interval4h: "4H",
types.Interval6h: "6H",
types.Interval12h: "12H",
types.Interval1d: "1D",
types.Interval2d: "2D",
types.Interval3d: "3D",
types.Interval5d: "5D",
types.Interval1w: "1W",
types.Interval1mo: "1M",
types.Interval3mo: "3M",
}
// kline // kline
type OkxExchange struct { type OkxExchange struct {
conf config.OkxExchange conf config.OkxExchange

92
internal/exchange/okx/okx_fetch.go

@ -0,0 +1,92 @@
package okx
import (
"context"
"errors"
"fmt"
"net/http"
"net/url"
"sig-pub/pkg/types"
"strings"
"time"
"github.com/go-resty/resty/v2"
"golang.org/x/time/rate"
)
const (
HttpBaseUrl = "https://www.okx.com"
KlineBefore0 int64 = 1672502400000 // k线开始数据 2023-01-01 00:00:00 GMT+8
)
type OkxFetcher struct {
client *resty.Client
httpProxy string
historyKlineLimiter *rate.Limiter
}
func NewOkxFetcher(httpProxy string) (f *OkxFetcher) {
client := resty.New()
client.SetTimeout(30 * time.Second)
client.SetTransport(&http.Transport{
MaxIdleConns: 100,
MaxConnsPerHost: 10,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 10 * time.Second,
})
if httpProxy != "" {
client.SetProxy(httpProxy)
}
f = &OkxFetcher{
client: client,
httpProxy: httpProxy,
historyKlineLimiter: rate.NewLimiter(rate.Every(100*time.Millisecond), 20), // rate: 20次/2s
}
return
}
// FetchHistoryKlines 获取交易产品历史K线数据
// https://my.okx.com/docs-v5/zh/#order-book-trading-market-data-get-candlesticks-history
// 周期区间 after > before, (after, before)
func (f *OkxFetcher) FetchHistoryKlines(ctx context.Context, okxInstId string, interval types.Interval, after, before int64) (klines []*types.Kline, err error) {
if okxInstId == "" {
err = errors.New("instid is empty")
return
}
if after <= 0 && before <= 0 {
err = errors.New("time range zero")
return
}
if err = f.historyKlineLimiter.Wait(ctx); err != nil {
return
}
var params []string
params = append(params, fmt.Sprintf("instId=%s", url.QueryEscape(okxInstId)))
params = append(params, "limit=100") // 最大为100
if v, ok := subscribeCandles[interval]; ok {
params = append(params, "bar="+v)
}
if after > 0 {
params = append(params, fmt.Sprintf("after=%d", after))
}
if before > 0 {
params = append(params, fmt.Sprintf("before=%d", before))
}
url := fmt.Sprintf("%s/api/v5/market/history-candles?%s", HttpBaseUrl, strings.Join(params, "&"))
resp, err := f.client.R().Get(url)
if err != nil {
return
}
status := resp.StatusCode()
if status != 200 {
err = fmt.Errorf("request history klines status error: %s, %s", url, resp.Status())
return
}
fmt.Println(string(resp.Body()))
return
}

44
internal/exchange/okx/okx_fetch_test.go

@ -0,0 +1,44 @@
package okx
import (
"context"
"fmt"
"sig-pub/pkg/types"
"testing"
"time"
)
func TestFetchHistoryKlines(t *testing.T) {
okxFetcher := NewOkxFetcher("http://192.168.1.5:7890")
// get inst+interval before, if=0 -> global before
// for interval, adder := range types.SupportedIntervals {
// _, _ = interval, adder
// }
interval := types.Interval5m
intervalAdder := types.SupportedIntervals[interval]
before := intervalAdder(KlineBefore0, -1)
for range 10 {
after := intervalAdder(before, 10)
klines, err := okxFetcher.FetchHistoryKlines(context.Background(), "BTC-USDT", interval, after, before)
if err != nil {
t.Error(err)
return
}
lastTs := intervalAdder(after, -1) // todo ts(last kline)-1
before = lastTs
for _, kline := range klines {
fmt.Println(kline)
}
fmt.Println("-----------------------------------------------------")
}
}
func TestA(t *testing.T) {
begin := time.UnixMilli(KlineBefore0)
before := begin.AddDate(0, -1, 0)
after := begin.AddDate(0, 3, 0)
fmt.Println("before:", before.UnixMilli())
fmt.Println("after:", after.UnixMilli())
}

12
pkg/storage/tsdb/victoria_metrics/metric.go

@ -55,12 +55,12 @@ func Kline2Metrics(inst types.TradeInstance, klines []*types.Kline) (metrics []*
ms, ok := instMetrics[inst.InstId] ms, ok := instMetrics[inst.InstId]
if !ok || (rawValueLimit > 0 && len(ms[0].Values) >= rawValueLimit) { if !ok || (rawValueLimit > 0 && len(ms[0].Values) >= rawValueLimit) {
ms = [6]*Metric{ ms = [6]*Metric{
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "open"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "open"),
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "high"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "high"),
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "low"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "low"),
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "close"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "close"),
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "vol"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "vol"),
NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "volQuote"), NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "volQuote"),
} }
instMetrics[inst.InstId] = ms instMetrics[inst.InstId] = ms
for i := range len(ms) { for i := range len(ms) {

68
pkg/storage/tsdb/victoria_metrics/vm.go

@ -8,10 +8,13 @@ import (
"io" "io"
"net/http" "net/http"
"net/url" "net/url"
"runtime/debug"
"sig-pub/pkg/config" "sig-pub/pkg/config"
"sig-pub/pkg/types" "sig-pub/pkg/types"
"sig-pub/pkg/zlog" "sig-pub/pkg/zlog"
"github.com/bytedance/sonic"
"github.com/govalues/decimal"
"github.com/klauspost/compress/zstd" "github.com/klauspost/compress/zstd"
) )
@ -81,6 +84,13 @@ func compressData(data []byte) ([]byte, error) {
// 获取原始k线列表 // 获取原始k线列表
func (vm *VictoriaMetricsTSDB) GetRangeKline(inst types.TradeInstance, interval types.Interval, start, end int64) (err error) { func (vm *VictoriaMetricsTSDB) GetRangeKline(inst types.TradeInstance, interval types.Interval, start, end int64) (err error) {
defer func() {
if r := recover(); r != nil {
zlog.Error("vmdb get range kline recover error:", r)
debug.PrintStack()
err = fmt.Errorf("%v", r)
}
}()
if inst.InstId == "" { if inst.InstId == "" {
err = errors.New("instid is empty") err = errors.New("instid is empty")
return return
@ -104,17 +114,65 @@ func (vm *VictoriaMetricsTSDB) GetRangeKline(inst types.TradeInstance, interval
// read response json line // read response json line
reader := bufio.NewReader(resp.Body) reader := bufio.NewReader(resp.Body)
var klines []*types.Kline
var vmLineBytes []byte
for { for {
line, err2 := reader.ReadBytes('\n') vmLineBytes, err = reader.ReadBytes('\n')
if err2 == io.EOF { if err == io.EOF || len(vmLineBytes) == 0 {
break break
} }
if err2 != nil { if err != nil {
zlog.Error(err) zlog.Error(err)
err = err2
return return
} }
zlog.Infof("response body line: %s", string(line))
vmMetric := new(VMMetricKline)
if err = sonic.Unmarshal(vmLineBytes, vmMetric); err != nil {
return
}
for i, ts := range vmMetric.Timestamps {
if len(klines) <= i {
klines = append(klines, &types.Kline{
Interval: types.Interval(vmMetric.Metric.Interval),
Ts: ts,
Confirm: true,
})
}
kline := klines[i]
if kline.Ts != ts {
err = errors.New("vm metric kline integrate error")
return
}
value := vmMetric.Values[i]
switch vmMetric.Metric.Kind {
case "open":
kline.Open = value
case "close":
kline.Close = value
case "high":
kline.High = value
case "low":
kline.Low = value
case "vol":
kline.Vol = value
case "volQuote":
kline.VolQuote = value
}
}
// zlog.Infof("response body line: %#v", vmMetric)
}
for _, kline := range klines {
zlog.Infof("kline: %#v", kline)
} }
return return
} }
type VMMetricKline struct {
Metric struct {
Name string `json:"__name__"`
Interval string `json:"interval"`
Kind string `json:"kind"`
} `json:"metric"`
Values []decimal.Decimal `json:"values"`
Timestamps []int64 `json:"timestamps"`
}

2
pkg/storage/tsdb/victoria_metrics/vm_test.go

@ -13,8 +13,8 @@ func TestGetRangeKline(t *testing.T) {
vmdb := NewVictoriaMetricsTSDB(config.VictoriaMetricsConfig{Addr: test_addr}) vmdb := NewVictoriaMetricsTSDB(config.VictoriaMetricsConfig{Addr: test_addr})
inst := types.TradeInstance{ inst := types.TradeInstance{
InstId: "BTC_USDT", InstId: "BTC_USDT",
Exchange: types.ExchangeOKX,
} }
vmdb.GetRangeKline(inst, types.Interval1m, 1753368047691, time.Now().UnixMilli()) vmdb.GetRangeKline(inst, types.Interval1m, 1753368047691, time.Now().UnixMilli())
} }

8
pkg/types/instance.go

@ -2,8 +2,8 @@ package types
// TradeInstance 交易产品 // TradeInstance 交易产品
type TradeInstance struct { type TradeInstance struct {
InstId string InstId string // 交易产品系统id
TickSz int32 Exchange Exchange // 当前处理交易产品交易所
MinSz int32 PriceSz int32 // 价格精度
// InstType string QuantitySz int32 // 交易量精度
} }

89
pkg/types/interval.go

@ -1,33 +1,44 @@
package types package types
import "time"
var LossEmoji = "🔥" var LossEmoji = "🔥"
var ProfitEmoji = "💰" var ProfitEmoji = "💰"
type Interval string type Interval string
func (i Interval) Minutes() (int64, bool) { // AddMul 对指定毫秒时间戳增加周期数
m, ok := SupportedIntervals[i] func (i Interval) AddMul(ts, mul int64) (int64, bool) {
if !ok || m <= 0 { c, ok := SupportedIntervals[i]
return m, false if !ok {
return ts, false
} }
return m / 60, true return c(ts, mul), true
} }
func (i Interval) Seconds() (int64, bool) { // func (i Interval) Minutes() (int64, bool) {
m, ok := SupportedIntervals[i] // c, ok := SupportedIntervals[i]
if !ok || m <= 0 { // if !ok || c <= 0 {
return m, false // return c, false
} // }
return m, true // return c / 60, true
} // }
func (i Interval) Milliseconds() (int64, bool) { // func (i Interval) Seconds() (int64, bool) {
m, ok := SupportedIntervals[i] // m, ok := SupportedIntervals[i]
if !ok || m <= 0 { // if !ok || m <= 0 {
return m, false // return m, false
} // }
return m * 1000, true // return m, true
} // }
// func (i Interval) Milliseconds() (int64, bool) {
// m, ok := SupportedIntervals[i]
// if !ok || m <= 0 {
// return m, false
// }
// return m * 1000, true
// }
var ( var (
Interval1s = Interval("1s") Interval1s = Interval("1s")
@ -62,25 +73,27 @@ type IntervalWindow struct {
RightWindow *int `json:"rightWindow"` RightWindow *int `json:"rightWindow"`
} }
type IntervalMap map[Interval]int64 type IntervalMap map[Interval]IntervalAdder
type IntervalAdder func(ts, mul int64) (ret int64)
var SupportedIntervals = IntervalMap{ var SupportedIntervals = IntervalMap{
Interval1s: 1, // Interval1s: func(ts, mul int64) (ret int64) { return ts + (1000 * mul) },
Interval1m: 1 * 60, Interval1m: func(ts, mul int64) (ret int64) { return ts + (1 * 60 * 1000 * mul) },
Interval3m: 3 * 60, Interval3m: func(ts, mul int64) (ret int64) { return ts + (3 * 60 * 1000 * mul) },
Interval5m: 5 * 60, Interval5m: func(ts, mul int64) (ret int64) { return ts + (5 * 60 * 1000 * mul) },
Interval15m: 15 * 60, Interval15m: func(ts, mul int64) (ret int64) { return ts + (15 * 60 * 1000 * mul) },
Interval30m: 30 * 60, Interval30m: func(ts, mul int64) (ret int64) { return ts + (30 * 60 * 1000 * mul) },
Interval1h: 60 * 60, Interval1h: func(ts, mul int64) (ret int64) { return ts + (60 * 60 * 1000 * mul) },
Interval2h: 60 * 60 * 2, Interval2h: func(ts, mul int64) (ret int64) { return ts + (2 * 60 * 60 * 1000 * mul) },
Interval4h: 60 * 60 * 4, Interval4h: func(ts, mul int64) (ret int64) { return ts + (4 * 60 * 60 * 1000 * mul) },
Interval6h: 60 * 60 * 6, Interval6h: func(ts, mul int64) (ret int64) { return ts + (4 * 60 * 60 * 1000 * mul) },
Interval12h: 60 * 60 * 12, Interval12h: func(ts, mul int64) (ret int64) { return ts + (12 * 60 * 60 * 1000 * mul) },
Interval1d: 60 * 60 * 24, Interval1d: func(ts, mul int64) (ret int64) { return ts + (24 * 60 * 60 * 1000 * mul) },
Interval2d: 60 * 60 * 24 * 2, Interval2d: func(ts, mul int64) (ret int64) { return ts + (2 * 24 * 60 * 60 * 1000 * mul) },
Interval3d: 60 * 60 * 24 * 3, Interval3d: func(ts, mul int64) (ret int64) { return ts + (3 * 24 * 60 * 60 * 1000 * mul) },
Interval5d: 60 * 60 * 24 * 5, Interval5d: func(ts, mul int64) (ret int64) { return ts + (5 * 24 * 60 * 60 * 1000 * mul) },
Interval1w: 60 * 60 * 24 * 7, Interval1w: func(ts, mul int64) (ret int64) { return ts + (7 * 24 * 60 * 60 * 1000 * mul) },
// Interval1mo: 60 * 60 * 24 * 30, Interval1mo: func(ts, mul int64) (ret int64) { return time.UnixMilli(ts).AddDate(0, int(mul), 0).UnixMilli() },
// Interval3mo: 60 * 60 * 24 * 30 * 3, Interval3mo: func(ts, mul int64) (ret int64) { return time.UnixMilli(ts).AddDate(0, int(3*mul), 0).UnixMilli() },
} }

Loading…
Cancel
Save