You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 

274 lines
6.6 KiB

package vmts
import (
"bufio"
"bytes"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"runtime/debug"
"sig-pub/pkg/config"
"sig-pub/pkg/types"
"sig-pub/pkg/utils/collect"
"sig-pub/pkg/utils/conver"
"sig-pub/pkg/zlog"
"strconv"
"github.com/bytedance/sonic"
"github.com/govalues/decimal"
"github.com/klauspost/compress/zstd"
)
// VictoriaMetricsTSDB 时序库
type VictoriaMetricsTSDB struct {
addr string
}
func NewVictoriaMetricsTSDB(conf config.VictoriaMetricsConfig) *VictoriaMetricsTSDB {
return &VictoriaMetricsTSDB{
addr: conf.Addr,
}
}
func (vm *VictoriaMetricsTSDB) SaveKlines(inst types.TradeInstance, klines []*types.Kline) (err error) {
metrics, err := Kline2Metrics(inst, klines)
if err != nil {
return
}
err = vm.batchWriteMetrics(metrics)
return
}
func (vm *VictoriaMetricsTSDB) batchWriteMetrics(metrics []*Metric) (err error) {
var buf bytes.Buffer
// gz := gzip.NewWriter(&buf)
var data []byte
for _, metric := range metrics {
data, err = metric.ToRowJson()
if err != nil {
return
}
buf.Write(data)
buf.Write([]byte("\n"))
// gz.Write(data)
// gz.Write([]byte("\n"))
}
// gz.Close()
// file, _ := os.Open(fmt.Sprintf("%d.json", time.Now().Unix()))
// defer file.Close()
// err = os.WriteFile(fmt.Sprintf("./%d.json", time.Now().Unix()), buf.Bytes(), os.ModeAppend)
// if err != nil {
// return
// }
// datas, err := compressData(buf.Bytes())
// {kind="high"}[30m]
resp, err := http.Post(fmt.Sprintf("%s/api/v1/import", vm.addr), "application/json", &buf)
if err != nil {
return
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
zlog.Info("batch import status error: ", resp.Status)
}
return
}
func compressData(data []byte) ([]byte, error) {
var b bytes.Buffer
encoder, _ := zstd.NewWriter(&b)
defer encoder.Close()
if _, err := encoder.Write(data); err != nil {
return nil, err
}
return b.Bytes(), nil
}
// 获取原始k线列表 [start, ..., end]
func (vm *VictoriaMetricsTSDB) ListRangeKline(inst types.TradeInstance, interval types.Interval, start, end int64) (klines []*types.Kline, 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
}
if start <= 0 || end <= 0 {
err = errors.New("invalid time range")
return
}
match := fmt.Sprintf("%s{exchange=\"%s\", interval=\"%s\"}", inst.InstId, inst.Exchange.String(), interval)
params := fmt.Sprintf("start=%d&end=%d&match[]=%s", start, end, url.QueryEscape(match))
resp, err := http.Get(fmt.Sprintf("%s/api/v1/export?%s", vm.addr, params))
if err != nil {
return
}
defer func() {
if err := resp.Body.Close(); err != nil {
zlog.Error("close vmtsdb response error:", err)
}
}()
var priceSz, quantitySz decimal.Decimal
if priceSz, err = decimal.Ten.PowInt(int(inst.PriceSz)); err != nil {
return
}
if quantitySz, err = decimal.Ten.PowInt(int(inst.QuantitySz)); err != nil {
return
}
// read response json line
reader := bufio.NewReader(resp.Body)
var vmLineBytes []byte
for {
// todo 优化
vmLineBytes, err = reader.ReadBytes('\n')
if err == io.EOF || len(vmLineBytes) == 0 {
err = nil
break
}
if err != nil {
zlog.Error(err)
return
}
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 "vol", "volQuote":
value, err = value.Quo(quantitySz)
default:
value, err = value.Quo(priceSz)
}
if err != nil {
return
}
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
}
// 数据延时 强制刷盘
// https://www.victoriametrics.com.cn/docs/query/#latency
func (vm *VictoriaMetricsTSDB) ForceFlush() (err error) {
resp, err := http.Get(fmt.Sprintf("%s/internal/force_flush", vm.addr))
if err != nil {
return
}
resp.Body.Close()
return
}
// https://www.victoriametrics.com.cn/docs/query/#instant-query
func (vm *VictoriaMetricsTSDB) QueryInstant(time int64, step types.Interval) {
}
// QueryRange 查询结果 query_range 接口调用 [start, ..., end]
// https://www.victoriametrics.com.cn/docs/query/#range-query
func (vm *VictoriaMetricsTSDB) QueryRange(start, end int64, step types.Interval, query string) (meticMatrixs []*types.MeticMatrix, err error) {
params := fmt.Sprintf("latency_offset=0&start=%d&end=%d&step=%s&query=%s", start, end, step, url.QueryEscape(query))
resp, err := http.Get(fmt.Sprintf("%s/api/v1/query_range?%s", vm.addr, params))
if err != nil {
return
}
defer func() {
if err := resp.Body.Close(); err != nil {
zlog.Error("close vmtsdb query_range response error:", err)
}
}()
bytes, err := io.ReadAll(resp.Body)
if err != nil {
return
}
var r = new(ResponseQuery)
if err = sonic.Unmarshal(bytes, r); err != nil {
return
}
if r.Status != QueryStatusSuccess {
err = fmt.Errorf("query_range response error: errorType=%s, error=%s", r.ErrorType, r.Error)
return
}
if collect.NotIn(r.Data.ResultType, "vector", "matrix") {
err = fmt.Errorf("query_range ResultType error: %s", r.Data.ResultType)
return
}
for _, result := range r.Data.Result {
m := &types.MeticMatrix{
Name: result.Metric["__name__"],
Exchange: result.Metric["exchange"],
Interval: result.Metric["interval"],
Kind: result.Metric["kind"],
}
var values [][]any
if r.Data.ResultType == "vector" {
values = [][]any{result.Value}
}
if r.Data.ResultType == "matrix" {
values = result.Values
}
m.Timestamps = make([]int64, 0, len(values))
m.Values = make([]float64, 0, len(values))
for _, value := range values {
ts := int64(conver.ToFloat64(value[0]) * 1000)
v, e := strconv.ParseFloat(value[1].(string), 64)
if e != nil {
err = e
return
}
m.Timestamps = append(m.Timestamps, ts)
m.Values = append(m.Values, v)
}
meticMatrixs = append(meticMatrixs, m)
}
return
}