From 5d802be0fd19488bbf0a417f842185b76046b452 Mon Sep 17 00:00:00 2001 From: strange Date: Mon, 28 Jul 2025 08:37:54 +0800 Subject: [PATCH] okx kline fetch --- cmd/test/test.go | 10 +-- go.mod | 4 +- go.sum | 4 + internal/exchange/exchange_grpc_server.go | 34 +++++--- internal/exchange/okx/channel_kline.go | 33 ++++--- internal/exchange/okx/okx.go | 25 ++++++ internal/exchange/okx/okx_fetch.go | 92 ++++++++++++++++++++ internal/exchange/okx/okx_fetch_test.go | 44 ++++++++++ pkg/storage/tsdb/victoria_metrics/metric.go | 12 +-- pkg/storage/tsdb/victoria_metrics/vm.go | 68 +++++++++++++-- pkg/storage/tsdb/victoria_metrics/vm_test.go | 4 +- pkg/types/instance.go | 8 +- pkg/types/interval.go | 89 +++++++++++-------- 13 files changed, 339 insertions(+), 88 deletions(-) create mode 100644 internal/exchange/okx/okx_fetch.go create mode 100644 internal/exchange/okx/okx_fetch_test.go diff --git a/cmd/test/test.go b/cmd/test/test.go index a8d0810..f2ba235 100644 --- a/cmd/test/test.go +++ b/cmd/test/test.go @@ -1,8 +1,6 @@ package main import ( - "sig-pub/api/pb" - "sig-pub/pkg/types" "sig-pub/pkg/zlog" "github.com/VictoriaMetrics/metrics" @@ -24,10 +22,10 @@ type BTC struct { func testDecimal() { // zlog.Init() - var interval = "1m" - interval0 := types.Interval(interval) - sec, _ := interval0.Seconds() - zlog.Infof("hello...: %s, %s, %d", pb.Exchange_OKX.String(), pb.Exchange_BINANCE.String(), sec) + // var interval = "1m" + // interval0 := types.Interval(interval) + // sec, _ := interval0.Seconds() + // zlog.Infof("hello...: %s, %s, %d", pb.Exchange_OKX.String(), pb.Exchange_BINANCE.String()) // data := `{"price":"1.23456"}` // btc := new(BTC) diff --git a/go.mod b/go.mod index 7a3a42b..11e2cc9 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/dsnet/golib/unitconv v1.0.2 github.com/fanjindong/go-cache v0.0.6 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/govalues/decimal v0.1.36 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.uber.org/zap v1.27.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/protobuf v1.36.6 gopkg.in/natefinch/lumberjack.v2 v2.2.1 @@ -79,7 +82,6 @@ require ( go.uber.org/multierr v1.11.0 // indirect golang.org/x/arch v0.15.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/text v0.23.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20250303144028-a0af3efb3deb // indirect diff --git a/go.sum b/go.sum index 7645bc3..d52ba49 100644 --- a/go.sum +++ b/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/validator/v10 v10.26.0 h1:SP05Nqhjcvz81uJaRfEV0YBSSSGMc/iMaVtFbr3Sw2k= 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/go.mod h1:OXbVy3sEdcQ2Doequ6Z5BW6fXNQTmx+9S1MCJN5yJMI= 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.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY= 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-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= diff --git a/internal/exchange/exchange_grpc_server.go b/internal/exchange/exchange_grpc_server.go index e4a5ca1..9f3e990 100644 --- a/internal/exchange/exchange_grpc_server.go +++ b/internal/exchange/exchange_grpc_server.go @@ -97,6 +97,11 @@ func (svc *ExchangeGrpcServer) subscribeExchanges() { // consumerKline 消费交易所k线数据 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 { channelK, ok := <-c if !ok { @@ -115,18 +120,6 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types 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 pubMsgMap := make(map[string]*pb.StreamKline) @@ -138,11 +131,13 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types } exchangeName := exchange.String() + var confirmKlines []*types.Kline for _, kline := range channelK.Klines { // zlog.Infof("recv kline: %#v", kline) confirm := 0 if kline.Confirm { confirm = 1 + confirmKlines = append(confirmKlines, kline) } pubKey := fmt.Sprintf("/kline/%s/%s/%s/%d", exchangeName, exInst.InstId, kline.Interval, confirm) // todo 优化没有订阅者就跳过 @@ -158,6 +153,21 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types 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 { if len(msg.Klines) == 0 { continue diff --git a/internal/exchange/okx/channel_kline.go b/internal/exchange/okx/channel_kline.go index 01983d6..0387168 100644 --- a/internal/exchange/okx/channel_kline.go +++ b/internal/exchange/okx/channel_kline.go @@ -14,19 +14,20 @@ import ( type CandleData [][]string var ( - // k线类型 - // 如 [1s/1m/3m/5m/15m/30m/1H/2H/4H] - // 香港时间开盘价k线:[6H/12H/1D/2D/3D/1W/1M/3M] - // UTC时间开盘价k线:[6Hutc/12Hutc/1Dutc/2Dutc/3Dutc/1Wutc/1Mutc/3Mutc] - // "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"} - subCandles = []string{ - // "candle3M", "candle1M", - "candle1W", "candle1D", "candle2D", "candle3D", "candle5D", - "candle12H", "candle6H", "candle4H", "candle2H", "candle1H", - "candle30m", "candle15m", "candle5m", "candle3m", "candle1m", - "candle1s", - } +// k线类型 +// 如 [1s/1m/3m/5m/15m/30m/1H/2H/4H] +// 香港时间开盘价k线:[6H/12H/1D/2D/3D/1W/1M/3M] +// UTC时间开盘价k线:[6Hutc/12Hutc/1Dutc/2Dutc/3Dutc/1Wutc/1Mutc/3Mutc] +// "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"} +// +// subCandles = []string{ +// "candle3M", "candle1M", +// "candle1W", "candle1D", "candle2D", "candle3D", "candle5D", +// "candle12H", "candle6H", "candle4H", "candle2H", "candle1H", +// "candle30m", "candle15m", "candle5m", "candle3m", "candle1m", +// "candle1s", +// } ) // ChannelCandle k线订阅频道 @@ -39,11 +40,15 @@ func NewChannelCandle( traceId string, httpProxy string, ) *ChannelCandle { + var candles []string + for _, interval := range subscribeCandles { + candles = append(candles, "candle"+interval) + } cfg := wsChannelConfig[*CandleData, *types.ChannelKline]{ channelId: traceId, httpProxy: httpProxy, wsUrl: "/ws/v5/business", - subscribeChannels: subCandles, + subscribeChannels: candles, dataInstanceFunc: func() *CandleData { var d CandleData return &d diff --git a/internal/exchange/okx/okx.go b/internal/exchange/okx/okx.go index cbb5b49..d4f8261 100644 --- a/internal/exchange/okx/okx.go +++ b/internal/exchange/okx/okx.go @@ -5,6 +5,31 @@ import ( "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 type OkxExchange struct { conf config.OkxExchange diff --git a/internal/exchange/okx/okx_fetch.go b/internal/exchange/okx/okx_fetch.go new file mode 100644 index 0000000..9bb5d42 --- /dev/null +++ b/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 +} diff --git a/internal/exchange/okx/okx_fetch_test.go b/internal/exchange/okx/okx_fetch_test.go new file mode 100644 index 0000000..7aff070 --- /dev/null +++ b/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()) +} diff --git a/pkg/storage/tsdb/victoria_metrics/metric.go b/pkg/storage/tsdb/victoria_metrics/metric.go index c9e3220..b40e192 100644 --- a/pkg/storage/tsdb/victoria_metrics/metric.go +++ b/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] if !ok || (rawValueLimit > 0 && len(ms[0].Values) >= rawValueLimit) { ms = [6]*Metric{ - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "open"), - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "high"), - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "low"), - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "close"), - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "vol"), - NewMetric(inst.InstId, "interval", string(kline.Interval), "kind", "volQuote"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "open"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "high"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "low"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "close"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "vol"), + NewMetric(inst.InstId, "interval", string(kline.Interval), "exchange", string(inst.Exchange), "kind", "volQuote"), } instMetrics[inst.InstId] = ms for i := range len(ms) { diff --git a/pkg/storage/tsdb/victoria_metrics/vm.go b/pkg/storage/tsdb/victoria_metrics/vm.go index a2cb539..7ec2eb2 100644 --- a/pkg/storage/tsdb/victoria_metrics/vm.go +++ b/pkg/storage/tsdb/victoria_metrics/vm.go @@ -8,10 +8,13 @@ import ( "io" "net/http" "net/url" + "runtime/debug" "sig-pub/pkg/config" "sig-pub/pkg/types" "sig-pub/pkg/zlog" + "github.com/bytedance/sonic" + "github.com/govalues/decimal" "github.com/klauspost/compress/zstd" ) @@ -81,6 +84,13 @@ func compressData(data []byte) ([]byte, error) { // 获取原始k线列表 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 == "" { err = errors.New("instid is empty") return @@ -104,17 +114,65 @@ func (vm *VictoriaMetricsTSDB) GetRangeKline(inst types.TradeInstance, interval // read response json line reader := bufio.NewReader(resp.Body) + var klines []*types.Kline + var vmLineBytes []byte for { - line, err2 := reader.ReadBytes('\n') - if err2 == io.EOF { + vmLineBytes, err = reader.ReadBytes('\n') + if err == io.EOF || len(vmLineBytes) == 0 { break } - if err2 != nil { + if err != nil { zlog.Error(err) - err = err2 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 } + +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"` +} diff --git a/pkg/storage/tsdb/victoria_metrics/vm_test.go b/pkg/storage/tsdb/victoria_metrics/vm_test.go index 5d5adca..6f5b8b2 100644 --- a/pkg/storage/tsdb/victoria_metrics/vm_test.go +++ b/pkg/storage/tsdb/victoria_metrics/vm_test.go @@ -12,9 +12,9 @@ var test_addr = "http://127.0.0.1:8428" func TestGetRangeKline(t *testing.T) { vmdb := NewVictoriaMetricsTSDB(config.VictoriaMetricsConfig{Addr: test_addr}) inst := types.TradeInstance{ - InstId: "BTC_USDT", + InstId: "BTC_USDT", + Exchange: types.ExchangeOKX, } vmdb.GetRangeKline(inst, types.Interval1m, 1753368047691, time.Now().UnixMilli()) - } diff --git a/pkg/types/instance.go b/pkg/types/instance.go index 1ae9826..8ff07cc 100644 --- a/pkg/types/instance.go +++ b/pkg/types/instance.go @@ -2,8 +2,8 @@ package types // TradeInstance 交易产品 type TradeInstance struct { - InstId string - TickSz int32 - MinSz int32 - // InstType string + InstId string // 交易产品系统id + Exchange Exchange // 当前处理交易产品交易所 + PriceSz int32 // 价格精度 + QuantitySz int32 // 交易量精度 } diff --git a/pkg/types/interval.go b/pkg/types/interval.go index ec98b04..6695719 100644 --- a/pkg/types/interval.go +++ b/pkg/types/interval.go @@ -1,33 +1,44 @@ package types +import "time" + var LossEmoji = "🔥" var ProfitEmoji = "💰" type Interval string -func (i Interval) Minutes() (int64, bool) { - m, ok := SupportedIntervals[i] - if !ok || m <= 0 { - return m, false +// AddMul 对指定毫秒时间戳增加周期数 +func (i Interval) AddMul(ts, mul int64) (int64, bool) { + c, ok := SupportedIntervals[i] + if !ok { + return ts, false } - return m / 60, true + return c(ts, mul), true } -func (i Interval) Seconds() (int64, bool) { - m, ok := SupportedIntervals[i] - if !ok || m <= 0 { - return m, false - } - return m, true -} +// func (i Interval) Minutes() (int64, bool) { +// c, ok := SupportedIntervals[i] +// if !ok || c <= 0 { +// return c, false +// } +// return c / 60, true +// } -func (i Interval) Milliseconds() (int64, bool) { - m, ok := SupportedIntervals[i] - if !ok || m <= 0 { - return m, false - } - return m * 1000, true -} +// func (i Interval) Seconds() (int64, bool) { +// m, ok := SupportedIntervals[i] +// if !ok || m <= 0 { +// return m, false +// } +// 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 ( Interval1s = Interval("1s") @@ -62,25 +73,27 @@ type IntervalWindow struct { 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{ - Interval1s: 1, - Interval1m: 1 * 60, - Interval3m: 3 * 60, - Interval5m: 5 * 60, - Interval15m: 15 * 60, - Interval30m: 30 * 60, - Interval1h: 60 * 60, - Interval2h: 60 * 60 * 2, - Interval4h: 60 * 60 * 4, - Interval6h: 60 * 60 * 6, - Interval12h: 60 * 60 * 12, - Interval1d: 60 * 60 * 24, - Interval2d: 60 * 60 * 24 * 2, - Interval3d: 60 * 60 * 24 * 3, - Interval5d: 60 * 60 * 24 * 5, - Interval1w: 60 * 60 * 24 * 7, - // Interval1mo: 60 * 60 * 24 * 30, - // Interval3mo: 60 * 60 * 24 * 30 * 3, + // Interval1s: func(ts, mul int64) (ret int64) { return ts + (1000 * mul) }, + Interval1m: func(ts, mul int64) (ret int64) { return ts + (1 * 60 * 1000 * mul) }, + Interval3m: func(ts, mul int64) (ret int64) { return ts + (3 * 60 * 1000 * mul) }, + Interval5m: func(ts, mul int64) (ret int64) { return ts + (5 * 60 * 1000 * mul) }, + Interval15m: func(ts, mul int64) (ret int64) { return ts + (15 * 60 * 1000 * mul) }, + Interval30m: func(ts, mul int64) (ret int64) { return ts + (30 * 60 * 1000 * mul) }, + Interval1h: func(ts, mul int64) (ret int64) { return ts + (60 * 60 * 1000 * mul) }, + Interval2h: func(ts, mul int64) (ret int64) { return ts + (2 * 60 * 60 * 1000 * mul) }, + Interval4h: func(ts, mul int64) (ret int64) { return ts + (4 * 60 * 60 * 1000 * mul) }, + Interval6h: func(ts, mul int64) (ret int64) { return ts + (4 * 60 * 60 * 1000 * mul) }, + Interval12h: func(ts, mul int64) (ret int64) { return ts + (12 * 60 * 60 * 1000 * mul) }, + Interval1d: func(ts, mul int64) (ret int64) { return ts + (24 * 60 * 60 * 1000 * mul) }, + Interval2d: func(ts, mul int64) (ret int64) { return ts + (2 * 24 * 60 * 60 * 1000 * mul) }, + Interval3d: func(ts, mul int64) (ret int64) { return ts + (3 * 24 * 60 * 60 * 1000 * mul) }, + Interval5d: func(ts, mul int64) (ret int64) { return ts + (5 * 24 * 60 * 60 * 1000 * mul) }, + Interval1w: func(ts, mul int64) (ret int64) { return ts + (7 * 24 * 60 * 60 * 1000 * mul) }, + Interval1mo: func(ts, mul int64) (ret int64) { return time.UnixMilli(ts).AddDate(0, int(mul), 0).UnixMilli() }, + Interval3mo: func(ts, mul int64) (ret int64) { return time.UnixMilli(ts).AddDate(0, int(3*mul), 0).UnixMilli() }, }