Browse Source

initial history klines

main
strange 1 year ago
parent
commit
419d3dbea3
  1. 24
      cmd/exchange/main.go
  2. 34
      internal/exchange/exchange.go
  3. 244
      internal/exchange/exchange_grpc_server.go
  4. 8
      internal/exchange/okx/channel_kline.go
  5. 45
      internal/exchange/okx/okx_fetch.go
  6. 6
      internal/exchange/okx/okx_fetch_test.go
  7. 12
      internal/exchange/okx/okx_limiter.go
  8. 16
      internal/exchange/okx/okx_subscriber.go
  9. 7
      internal/exchange/okx/types.go
  10. 10
      pkg/aside/trade_instance_client.go
  11. 53
      pkg/storage/kvrocks/kvrocks.go
  12. 13
      pkg/types/exchange.go
  13. 4
      pkg/types/instance.go

24
cmd/exchange/main.go

@ -11,6 +11,7 @@ import (
"sig-pub/pkg/config" "sig-pub/pkg/config"
"sig-pub/pkg/grpc/discovery" "sig-pub/pkg/grpc/discovery"
"sig-pub/pkg/grpc/interceptor" "sig-pub/pkg/grpc/interceptor"
"sig-pub/pkg/storage/kvrocks"
vmts "sig-pub/pkg/storage/tsdb/victoria_metrics" vmts "sig-pub/pkg/storage/tsdb/victoria_metrics"
"sig-pub/pkg/utils/exit" "sig-pub/pkg/utils/exit"
"sig-pub/pkg/zlog" "sig-pub/pkg/zlog"
@ -38,11 +39,6 @@ func main() {
panic(err) panic(err)
} }
okxExchange := okx.NewOkxExchange(conf.Exchange.Okx)
if err := okxExchange.Init(); err != nil {
panic(err)
}
etcdClient, err := clientv3.New(conf.Etcd) etcdClient, err := clientv3.New(conf.Etcd)
if err != nil { if err != nil {
panic(err) panic(err)
@ -65,7 +61,23 @@ func main() {
} }
marketClient := pb.NewMarketClient(conn) marketClient := pb.NewMarketClient(conn)
tradeInstanceAside := aside.NewTradeInstanceAside(marketClient) tradeInstanceAside := aside.NewTradeInstanceAside(marketClient)
exchangeService := exchange.NewExchangeGrpcServer(tradeInstanceAside, tsdbService, okxExchange)
// kvrocks db
kvdb := kvrocks.NewKVRocksDB(conf.Database.Kvrocks)
if err := kvdb.Ping(); err != nil {
panic(err)
}
// okx exchange
okxSubscriber := okx.NewOkxSubscriber(conf.Exchange.Okx)
if err := okxSubscriber.Init(); err != nil {
panic(err)
}
okxFetcher := okx.NewOkxFetcher(conf.Exchange.Okx.HttpProxy)
okxExchange := exchange.NewExchange(okxFetcher, okxSubscriber)
// exhcange main service
exchangeService := exchange.NewExchangeGrpcServer(tradeInstanceAside, tsdbService, kvdb, okxExchange)
if err := exchangeService.Init(); err != nil { if err := exchangeService.Init(); err != nil {
panic(err) panic(err)
} }

34
internal/exchange/exchange.go

@ -1,6 +1,11 @@
package exchange package exchange
import "sig-pub/pkg/types" import (
"context"
"fmt"
"sig-pub/pkg/types"
"sync"
)
// 交易所行情数据订阅 // 交易所行情数据订阅
type ExchangeSubscriber interface { type ExchangeSubscriber interface {
@ -19,4 +24,31 @@ type ExchangeSubscriber interface {
// 交易所行情数据请求 // 交易所行情数据请求
type ExchangeFetcher interface { type ExchangeFetcher interface {
// 交易所类型
ExhcangeType() types.Exchange
// 获取区间内历史k线数据
FetchHistoryKlines(ctx context.Context, instId string, interval types.Interval, after, before int64) (klines []*types.Kline, err error)
}
// 交易所交互接口
type Exchange struct {
ExType types.Exchange
Fetcher ExchangeFetcher
Subscriber ExchangeSubscriber
Insts map[string]*types.TradeInstance // <ExchangeInstId, Inst>
sync.RWMutex
}
func NewExchange(fetcher ExchangeFetcher, subscriber ExchangeSubscriber) *Exchange {
exType := fetcher.ExhcangeType()
if exType != subscriber.ExhcangeType() {
panic(fmt.Errorf("exchange type not match: %#v, %#v", exType, subscriber.ExhcangeType()))
}
return &Exchange{
ExType: exType,
Fetcher: fetcher,
Subscriber: subscriber,
Insts: make(map[string]*types.TradeInstance),
}
} }

244
internal/exchange/exchange_grpc_server.go

@ -6,7 +6,8 @@ import (
"io" "io"
"sig-pub/api/pb" "sig-pub/api/pb"
"sig-pub/pkg/aside" "sig-pub/pkg/aside"
"sig-pub/pkg/data/entity" "sig-pub/pkg/data"
"sig-pub/pkg/storage/kvrocks"
"sig-pub/pkg/types" "sig-pub/pkg/types"
"sig-pub/pkg/zlog" "sig-pub/pkg/zlog"
"sync" "sync"
@ -15,43 +16,34 @@ import (
"google.golang.org/grpc" "google.golang.org/grpc"
) )
type Exchange struct {
Type pb.Exchange
Subscriber ExchangeSubscriber
Insts map[string]*entity.TradeInstanceExchange
sync.RWMutex
}
type ExchangeGrpcServer struct { type ExchangeGrpcServer struct {
pb.UnimplementedExchangeServiceServer pb.UnimplementedExchangeServiceServer
exchangeMap map[pb.Exchange]*Exchange exchangeMap map[types.Exchange]*Exchange
tradeInstanceAside *aside.TradeInstanceAside tradeInstanceAside *aside.TradeInstanceAside
exchangeDataService *ExchangeDataService exchangeDataService *ExchangeDataService
kvdb *kvrocks.KVRocksDB
klineStreamId int64 klineStreamId int64
klinePublisher *Publisher[int64, grpc.BidiStreamingServer[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline]] klinePublisher *Publisher[int64, grpc.BidiStreamingServer[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline]]
} }
// exchanges: 支持的数据源交易所 // exchanges: 支持的数据源交易所
func NewExchangeGrpcServer(tradeInstanceAside *aside.TradeInstanceAside, exchangeDataService *ExchangeDataService, exchanges ...ExchangeSubscriber) *ExchangeGrpcServer { func NewExchangeGrpcServer(
exchangeMap := make(map[pb.Exchange]*Exchange) tradeInstanceAside *aside.TradeInstanceAside,
exchangeDataService *ExchangeDataService,
kvdb *kvrocks.KVRocksDB,
exchanges ...*Exchange,
) *ExchangeGrpcServer {
exchangeMap := make(map[types.Exchange]*Exchange)
for _, exchange := range exchanges { for _, exchange := range exchanges {
exchangeType := exchange.ExhcangeType() exchangeMap[exchange.ExType] = exchange
pbExchangeType, ok := exchangeType.Exchange2PB()
if !ok {
panic(fmt.Errorf("unknown exchange: %#v", exchangeType))
}
exchangeMap[pbExchangeType] = &Exchange{
Type: pbExchangeType,
Subscriber: exchange,
Insts: make(map[string]*entity.TradeInstanceExchange),
}
} }
return &ExchangeGrpcServer{ return &ExchangeGrpcServer{
exchangeMap: exchangeMap, exchangeMap: exchangeMap,
tradeInstanceAside: tradeInstanceAside, tradeInstanceAside: tradeInstanceAside,
exchangeDataService: exchangeDataService, exchangeDataService: exchangeDataService,
kvdb: kvdb,
klinePublisher: NewPublisher[int64, grpc.BidiStreamingServer[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline]](16), klinePublisher: NewPublisher[int64, grpc.BidiStreamingServer[pb.ReqStreamSubscribeKline, pb.RspStreamSubscribeKline]](16),
} }
} }
@ -68,16 +60,30 @@ func (svc *ExchangeGrpcServer) subscribeExchanges() {
for _, exchange := range svc.exchangeMap { for _, exchange := range svc.exchangeMap {
go func(exchange *Exchange) { go func(exchange *Exchange) {
// get exchange trade instances // get exchange trade instances
insts, err := svc.tradeInstanceAside.ListExchangeTradeInstance(context.Background(), exchange.Type) insts, err := svc.tradeInstanceAside.ListExchangeTradeInstance(context.Background(), exchange.ExType)
if err != nil { if err != nil {
zlog.Error(err) zlog.Error(err)
return return
} }
var exchangeInstIds []string var exchangeInstIds []string
var processingInsts []types.TradeInstance
exchange.Lock() exchange.Lock()
for _, inst := range insts { for _, inst := range insts {
exchangeInstIds = append(exchangeInstIds, inst.ExchangeInstId) exchangeInstIds = append(exchangeInstIds, inst.ExchangeInstId)
exchange.Insts[inst.ExchangeInstId] = inst tradeInst := &types.TradeInstance{
InstId: inst.InstId,
Status: inst.Status,
PriceSz: 0,
QuantitySz: 0,
ExchangeInstId: inst.ExchangeInstId,
Exchange: exchange.ExType,
}
exchange.Insts[inst.ExchangeInstId] = tradeInst
// 待初始化币种数据
if inst.Status == data.StatusProcessing {
processingInsts = append(processingInsts, *tradeInst)
}
} }
exchange.Unlock() exchange.Unlock()
@ -87,20 +93,23 @@ func (svc *ExchangeGrpcServer) subscribeExchanges() {
zlog.Error(err) zlog.Error(err)
return return
} }
go func() {
c := exchange.Subscriber.ConsumerKline() c := exchange.Subscriber.ConsumerKline()
svc.consumerKline(exchange, c) svc.consumerKline(exchange, c)
// todo subscribe books 订单簿 // todo subscribe books 订单簿
zlog.Infof("unsubscribe exchange: %s", exchange.Type.String()) zlog.Infof("unsubscribe exchange: %s", exchange.ExType)
}()
// 初始化k线数据
go svc.initialKlines(exchange, processingInsts)
}(exchange) }(exchange)
} }
} }
// 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) exchangeType := exchange.ExType
if !ok {
panic(fmt.Errorf("unknown exchange type: %v", exchange.Type))
}
for { for {
channelK, ok := <-c channelK, ok := <-c
@ -109,7 +118,7 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
} }
// 交易所 instid 转 sig-instid // 交易所 instid 转 sig-instid
var exInst *entity.TradeInstanceExchange var exInst *types.TradeInstance
exchange.RLock() exchange.RLock()
if inst, ok := exchange.Insts[channelK.InstId]; ok && inst != nil { if inst, ok := exchange.Insts[channelK.InstId]; ok && inst != nil {
exchange.RUnlock() exchange.RUnlock()
@ -124,12 +133,11 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
pubMsgMap := make(map[string]*pb.StreamKline) pubMsgMap := make(map[string]*pb.StreamKline)
// instId := channelK.InstId // instId := channelK.InstId
exchange, ok := channelK.Exchange.Exchange2PB() pbExType, err := channelK.Exchange.Exchange2PB()
if !ok { if err != nil {
zlog.Errorf("unknown exchange kline: %v", channelK.Exchange) zlog.Error(err)
continue continue
} }
exchangeName := exchange.String()
var confirmKlines []*types.Kline var confirmKlines []*types.Kline
for _, kline := range channelK.Klines { for _, kline := range channelK.Klines {
@ -139,13 +147,13 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
confirm = 1 confirm = 1
confirmKlines = append(confirmKlines, kline) 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", exchangeType, exInst.InstId, kline.Interval, confirm)
// todo 优化没有订阅者就跳过 // todo 优化没有订阅者就跳过
msg, ok := pubMsgMap[pubKey] msg, ok := pubMsgMap[pubKey]
if !ok { if !ok {
msg = new(pb.StreamKline) msg = new(pb.StreamKline)
msg.InstId = exInst.InstId msg.InstId = exInst.InstId
msg.Exchange = exchange msg.Exchange = pbExType
pubMsgMap[pubKey] = msg pubMsgMap[pubKey] = msg
} }
@ -155,14 +163,8 @@ func (svc *ExchangeGrpcServer) consumerKline(exchange *Exchange, c <-chan *types
// tsdb storage // tsdb storage
if len(confirmKlines) > 0 { if len(confirmKlines) > 0 {
typeInst := types.TradeInstance{
InstId: exInst.InstId, // channelK.InstId
PriceSz: 0,
QuantitySz: 0,
Exchange: exchangeType,
}
// todo 异步处理 // todo 异步处理
err := svc.exchangeDataService.SaveKlines(typeInst, confirmKlines) err := svc.exchangeDataService.SaveKlines(*exInst, confirmKlines)
if err != nil { if err != nil {
zlog.Errorf("kline save to tsdb error: ", err) zlog.Errorf("kline save to tsdb error: ", err)
} }
@ -279,3 +281,161 @@ func (svc *ExchangeGrpcServer) SubscribeKline(stream grpc.BidiStreamingServer[pb
// } // }
// } // }
} }
const (
KlineBefore0 int64 = 1672502400000 // k线开始数据 2023-01-01 00:00:00 GMT+8
HistoryKlineTsKey string = "history-kline-ts:%s/%s/%s" // exchange:sig-instid:interval
)
// initialKline 初始化交易产品历史k线数据
func (svc *ExchangeGrpcServer) initialKlines(exchange *Exchange, insts []types.TradeInstance) {
var err error
defer func() {
if err != nil {
zlog.Error(err)
}
}()
type task struct {
inst types.TradeInstance
interval types.Interval
afterTs int64
beforeTs int64
}
// 任务 channel
ch := make(chan task)
// 任务生成器
go func() {
for _, inst := range insts {
// for interval, intervalAdder := range types.SupportedIntervals {
interval := types.Interval1h
intervalAdder := types.SupportedIntervals[interval]
tsKey := fmt.Sprintf(HistoryKlineTsKey, inst.Exchange, inst.InstId, interval)
beforeTs, e := svc.kvdb.GetI64(context.Background(), tsKey)
if e != nil {
err = e
return
}
if beforeTs == 0 {
beforeTs = intervalAdder(KlineBefore0, -1)
}
for {
afterTs := intervalAdder(beforeTs, 101)
ch <- task{
inst: inst,
interval: interval,
afterTs: afterTs,
beforeTs: beforeTs,
}
beforeTs = afterTs
// todo compare ws live ts
}
// }
}
close(ch)
}()
// 任务消费器 8协程并行
wg := new(sync.WaitGroup)
for range 11 {
wg.Add(1)
go func() {
defer wg.Done()
for {
task, ok := <-ch
if !ok {
break
}
interval, afterTs, beforeTs := task.interval, task.afterTs, task.beforeTs
klines, e := exchange.Fetcher.FetchHistoryKlines(context.Background(), task.inst.ExchangeInstId, interval, afterTs, beforeTs)
if e != nil {
// todo retry
err = e
return
}
if len(klines) == 0 {
continue // todo ....
}
zlog.Infof("fetch interval %s %d~%d klines: ret=%d~%d, %d klines", interval, afterTs, beforeTs, klines[0].Ts, klines[len(klines)-1].Ts, len(klines))
// store to tsdb
err = svc.exchangeDataService.SaveKlines(task.inst, klines)
if err != nil {
return
}
}
}()
}
wg.Wait()
zlog.Infof("%d insts initial finished", len(insts))
// var intervals []types.Interval
// for interval := range types.SupportedIntervals {
// intervals = append(intervals, interval)
// }
// var intervalsTs = make([]int, len(intervals))
// var index int
// var lock sync.Mutex
// var getTask = func() (interval types.Interval, afterTs, beforeTs int64) {
// lock.Lock()
// ts := intervalsTs[index]
// if ts == -1 {
// index++
// }
// lock.Unlock()
// return
// }
// var finishTask = func(interval types.Interval) {
// }
// ctx := context.Background()
// inst := insts[0]
// // for interval, intervalAdder := range types.SupportedIntervals {
// interval := types.Interval1h
// intervalAdder := types.SupportedIntervals[interval]
// // kvrocks get exchange+inst+interval last/ts
// tsKey := fmt.Sprintf(HistoryKlineTsKey, inst.Exchange, inst.InstId, interval)
// beforeTs, e := svc.kvdb.GetI64(ctx, tsKey)
// if e != nil {
// err = e
// return
// }
// if beforeTs == 0 {
// beforeTs = intervalAdder(KlineBefore0, -1)
// }
// for {
// afterTs := intervalAdder(beforeTs, 101)
// klines, e := exchange.Fetcher.FetchHistoryKlines(ctx, inst.ExchangeInstId, interval, afterTs, beforeTs)
// if e != nil {
// err = e
// return
// }
// if len(klines) == 0 {
// break
// }
// zlog.Infof("fetch interval %s %d~%d klines: ret=%d~%d, %d klines", interval, afterTs, beforeTs, klines[0].Ts, klines[len(klines)-1].Ts, len(klines))
// // store to tsdb
// err = svc.exchangeDataService.SaveKlines(inst, klines)
// if err != nil {
// return
// }
// // set kvdb inst ts mark
// beforeTs = klines[0].Ts
// }
// // }
// zlog.Infof("%s %s initial finished", inst.Exchange, inst.InstId)
}

8
internal/exchange/okx/channel_kline.go

@ -139,10 +139,10 @@ func candleData2Klines(channelData *ChannelData[*CandleData]) (r *types.ChannelK
kline.Interval = types.Interval5d kline.Interval = types.Interval5d
case "candle1W": case "candle1W":
kline.Interval = types.Interval1w kline.Interval = types.Interval1w
// case "candle1M": case "candle1M":
// kline.Interval = types.Interval1mo kline.Interval = types.Interval1mo
// case "candle3M": case "candle3M":
// kline.Interval = types.Interval3mo kline.Interval = types.Interval3mo
} }
var kinds = []*decimal.Decimal{&kline.Open, &kline.High, &kline.Low, &kline.Close, &kline.Vol, nil, &kline.VolQuote} var kinds = []*decimal.Decimal{&kline.Open, &kline.High, &kline.Low, &kline.Close, &kline.Vol, nil, &kline.VolQuote}
for i := 1; i <= 7; i++ { for i := 1; i <= 7; i++ {

45
internal/exchange/okx/okx_fetch.go

@ -7,22 +7,22 @@ import (
"net/http" "net/http"
"net/url" "net/url"
"sig-pub/pkg/types" "sig-pub/pkg/types"
"strconv"
"strings" "strings"
"time" "time"
"github.com/bytedance/sonic"
"github.com/go-resty/resty/v2" "github.com/go-resty/resty/v2"
"golang.org/x/time/rate" "github.com/govalues/decimal"
) )
const ( const (
HttpBaseUrl = "https://www.okx.com" HttpBaseUrl = "https://www.okx.com"
KlineBefore0 int64 = 1672502400000 // k线开始数据 2023-01-01 00:00:00 GMT+8
) )
type OkxFetcher struct { type OkxFetcher struct {
client *resty.Client client *resty.Client
httpProxy string httpProxy string
historyKlineLimiter *rate.Limiter
} }
func NewOkxFetcher(httpProxy string) (f *OkxFetcher) { func NewOkxFetcher(httpProxy string) (f *OkxFetcher) {
@ -41,11 +41,14 @@ func NewOkxFetcher(httpProxy string) (f *OkxFetcher) {
f = &OkxFetcher{ f = &OkxFetcher{
client: client, client: client,
httpProxy: httpProxy, httpProxy: httpProxy,
historyKlineLimiter: rate.NewLimiter(rate.Every(100*time.Millisecond), 20), // rate: 20次/2s
} }
return return
} }
func (okx *OkxFetcher) ExhcangeType() types.Exchange {
return types.ExchangeOKX
}
// FetchHistoryKlines 获取交易产品历史K线数据 // FetchHistoryKlines 获取交易产品历史K线数据
// https://my.okx.com/docs-v5/zh/#order-book-trading-market-data-get-candlesticks-history // https://my.okx.com/docs-v5/zh/#order-book-trading-market-data-get-candlesticks-history
// 周期区间 after > before, (after, before) // 周期区间 after > before, (after, before)
@ -58,7 +61,7 @@ func (f *OkxFetcher) FetchHistoryKlines(ctx context.Context, okxInstId string, i
err = errors.New("time range zero") err = errors.New("time range zero")
return return
} }
if err = f.historyKlineLimiter.Wait(ctx); err != nil { if err = fetchHistoryKlineLimiter.Wait(ctx); err != nil {
return return
} }
@ -87,6 +90,36 @@ func (f *OkxFetcher) FetchHistoryKlines(ctx context.Context, okxInstId string, i
return return
} }
fmt.Println(string(resp.Body())) r := new(RespHistoryKline)
if err = sonic.Unmarshal(resp.Body(), r); err != nil {
return
}
if r.Code != "0" {
err = fmt.Errorf("code: %s, msg: %s", r.Code, r.Msg)
return
}
for _, data := range r.Data {
ts, e := strconv.ParseInt(data[0], 10, 64)
if e != nil {
err = e
return
}
kline := &types.Kline{
Interval: interval,
Ts: ts,
Confirm: data[8] == "1",
}
var kinds = []*decimal.Decimal{&kline.Open, &kline.High, &kline.Low, &kline.Close, &kline.Vol, nil, &kline.VolQuote}
for i := 1; i <= 7; i++ {
if kinds[i-1] == nil {
continue
}
*kinds[i-1], err = decimal.Parse(data[i])
if err != nil {
return
}
}
klines = append(klines, kline)
}
return return
} }

6
internal/exchange/okx/okx_fetch_test.go

@ -8,6 +8,8 @@ import (
"time" "time"
) )
var klineBefore0 int64 = 1672502400000 // k线开始数据 2023-01-01 00:00:00 GMT+8
func TestFetchHistoryKlines(t *testing.T) { func TestFetchHistoryKlines(t *testing.T) {
okxFetcher := NewOkxFetcher("http://192.168.1.5:7890") okxFetcher := NewOkxFetcher("http://192.168.1.5:7890")
@ -18,7 +20,7 @@ func TestFetchHistoryKlines(t *testing.T) {
// } // }
interval := types.Interval5m interval := types.Interval5m
intervalAdder := types.SupportedIntervals[interval] intervalAdder := types.SupportedIntervals[interval]
before := intervalAdder(KlineBefore0, -1) before := intervalAdder(klineBefore0, -1)
for range 10 { for range 10 {
after := intervalAdder(before, 10) after := intervalAdder(before, 10)
klines, err := okxFetcher.FetchHistoryKlines(context.Background(), "BTC-USDT", interval, after, before) klines, err := okxFetcher.FetchHistoryKlines(context.Background(), "BTC-USDT", interval, after, before)
@ -36,7 +38,7 @@ func TestFetchHistoryKlines(t *testing.T) {
} }
func TestA(t *testing.T) { func TestA(t *testing.T) {
begin := time.UnixMilli(KlineBefore0) begin := time.UnixMilli(klineBefore0)
before := begin.AddDate(0, -1, 0) before := begin.AddDate(0, -1, 0)
after := begin.AddDate(0, 3, 0) after := begin.AddDate(0, 3, 0)
fmt.Println("before:", before.UnixMilli()) fmt.Println("before:", before.UnixMilli())

12
internal/exchange/okx/okx_limiter.go

@ -0,0 +1,12 @@
package okx
import (
"time"
"golang.org/x/time/rate"
)
var (
// https://my.okx.com/docs-v5/zh/#order-book-trading-market-data-get-candlesticks-history
fetchHistoryKlineLimiter = rate.NewLimiter(rate.Every(100*time.Millisecond), 2)
)

16
internal/exchange/okx/okx.go → internal/exchange/okx/okx_subscriber.go

@ -31,18 +31,18 @@ var subscribeCandles = map[types.Interval]string{
} }
// kline // kline
type OkxExchange struct { type OkxSubscriber struct {
conf config.OkxExchange conf config.OkxExchange
channelCandle *ChannelCandle // K线频道 channelCandle *ChannelCandle // K线频道
} }
func NewOkxExchange(conf config.OkxExchange) *OkxExchange { func NewOkxSubscriber(conf config.OkxExchange) *OkxSubscriber {
return &OkxExchange{ return &OkxSubscriber{
conf: conf, conf: conf,
} }
} }
func (okx *OkxExchange) Init() (err error) { func (okx *OkxSubscriber) Init() (err error) {
// TODO 多个 ChannelCandle 实例 OkxAggregate // TODO 多个 ChannelCandle 实例 OkxAggregate
okx.channelCandle = NewChannelCandle("candle-0", okx.conf.HttpProxy) okx.channelCandle = NewChannelCandle("candle-0", okx.conf.HttpProxy)
if err = okx.channelCandle.Init(); err != nil { if err = okx.channelCandle.Init(); err != nil {
@ -51,20 +51,20 @@ func (okx *OkxExchange) Init() (err error) {
return return
} }
func (okx *OkxExchange) ExhcangeType() types.Exchange { func (okx *OkxSubscriber) ExhcangeType() types.Exchange {
return types.ExchangeOKX return types.ExchangeOKX
} }
func (okx *OkxExchange) ConsumerKline() <-chan *types.ChannelKline { func (okx *OkxSubscriber) ConsumerKline() <-chan *types.ChannelKline {
return okx.channelCandle.Consumer() return okx.channelCandle.Consumer()
} }
// 订阅产品k线行情 // 订阅产品k线行情
func (okx *OkxExchange) SubscribeKline(instIds ...string) (err error) { func (okx *OkxSubscriber) SubscribeKline(instIds ...string) (err error) {
return okx.channelCandle.Subscribe(instIds...) return okx.channelCandle.Subscribe(instIds...)
} }
// 取消订阅产品k线行情 // 取消订阅产品k线行情
func (okx *OkxExchange) UnsubscribeKline(instIds ...string) (err error) { func (okx *OkxSubscriber) UnsubscribeKline(instIds ...string) (err error) {
return okx.channelCandle.Unsubscribe(instIds...) return okx.channelCandle.Unsubscribe(instIds...)
} }

7
internal/exchange/okx/types.go

@ -65,3 +65,10 @@ type MarketData struct {
Volume float64 Volume float64
Timestamp time.Time Timestamp time.Time
} }
// RespHistoryKline 历史k线响应
type RespHistoryKline struct {
Code string `json:"code"`
Msg string `json:"msg"`
Data [][]string `json:"data"`
}

10
pkg/aside/trade_instance_client.go

@ -5,6 +5,7 @@ import (
"sig-pub/api/pb" "sig-pub/api/pb"
"sig-pub/pkg/data/entity" "sig-pub/pkg/data/entity"
"sig-pub/pkg/mapping" "sig-pub/pkg/mapping"
"sig-pub/pkg/types"
"sig-pub/pkg/utils/kvcache" "sig-pub/pkg/utils/kvcache"
"time" "time"
@ -61,9 +62,14 @@ func (c *TradeInstanceAside) getTradeInstance0(ctx context.Context, instId strin
} }
// ListExchangeTradeInstance 获取交易所支持的交易实例 // ListExchangeTradeInstance 获取交易所支持的交易实例
func (c *TradeInstanceAside) ListExchangeTradeInstance(ctx context.Context, exchange pb.Exchange) (exInsts []*entity.TradeInstanceExchange, err error) { func (c *TradeInstanceAside) ListExchangeTradeInstance(ctx context.Context, exchange types.Exchange) (exInsts []*entity.TradeInstanceExchange, err error) {
pbExType, err := exchange.Exchange2PB()
if err != nil {
return
}
rsp, err := c.client.ListExchangeTradeInstance(ctx, &pb.ReqListExchangeTradeInstance{ rsp, err := c.client.ListExchangeTradeInstance(ctx, &pb.ReqListExchangeTradeInstance{
Exchanges: []pb.Exchange{exchange}, Exchanges: []pb.Exchange{pbExType},
}) })
if err != nil { if err != nil {
return return

53
pkg/storage/kvrocks/kvrocks.go

@ -0,0 +1,53 @@
package kvrocks
import (
"context"
"fmt"
"strconv"
"time"
"github.com/redis/go-redis/v9"
)
type KVRocksDB struct {
client *redis.Client
}
func NewKVRocksDB(conf redis.Options) *KVRocksDB {
client := redis.NewClient(&conf)
return &KVRocksDB{
client: client,
}
}
func (db *KVRocksDB) Ping() (err error) {
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
_, err = db.client.Ping(ctx).Result()
if err != nil {
return fmt.Errorf("can't connect kvrocks, %v", err)
}
return
}
func (db *KVRocksDB) DB() *redis.Client {
return db.client
}
func (db *KVRocksDB) GetI64(ctx context.Context, key string) (v int64, err error) {
v, err = db.DB().Get(ctx, key).Int64()
if err != nil {
if err == redis.Nil {
return 0, nil
}
return 0, err
}
return v, nil
}
// todo set async
func (db *KVRocksDB) SetI64(ctx context.Context, key string, v int64) (err error) {
vs := strconv.FormatInt(v, 10)
err = db.DB().Set(ctx, key, vs, 0).Err()
return
}

13
pkg/types/exchange.go

@ -1,6 +1,9 @@
package types package types
import "sig-pub/api/pb" import (
"fmt"
"sig-pub/api/pb"
)
type Exchange string type Exchange string
@ -10,14 +13,14 @@ var (
ExchangeBINANCE = Exchange(pb.Exchange_BINANCE.String()) // 币安 ExchangeBINANCE = Exchange(pb.Exchange_BINANCE.String()) // 币安
) )
func (ex Exchange) Exchange2PB() (pb.Exchange, bool) { func (ex Exchange) Exchange2PB() (pb.Exchange, error) {
switch ex { switch ex {
case ExchangeOKX: case ExchangeOKX:
return pb.Exchange_OKX, true return pb.Exchange_OKX, nil
case ExchangeBINANCE: case ExchangeBINANCE:
return pb.Exchange_BINANCE, true return pb.Exchange_BINANCE, nil
default: default:
return pb.Exchange_SIG, false return pb.Exchange_SIG, fmt.Errorf("unknown exchange: %v", ex)
} }
} }

4
pkg/types/instance.go

@ -3,7 +3,9 @@ package types
// TradeInstance 交易产品 // TradeInstance 交易产品
type TradeInstance struct { type TradeInstance struct {
InstId string // 交易产品系统id InstId string // 交易产品系统id
Exchange Exchange // 当前处理交易产品交易所
PriceSz int32 // 价格精度 PriceSz int32 // 价格精度
QuantitySz int32 // 交易量精度 QuantitySz int32 // 交易量精度
Status int32 // 交易所交易产品状态
ExchangeInstId string // 交易所交易产品id
Exchange Exchange // 当前处理交易产品交易所
} }

Loading…
Cancel
Save