Browse Source

initial fetch klines

main
strange 11 months ago
parent
commit
41a276b06e
  1. 18
      cmd/exchange/main.go
  2. 5
      cmd/market/main.go
  3. 4
      cmd/sig-admin/main.go
  4. 5
      docker-compose.yml
  5. 274
      internal/exchange/exchange_grpc_server.go
  6. 3
      pkg/grpc/discovery/consul_naming.go

18
cmd/exchange/main.go

@ -31,22 +31,9 @@ func main() {
conf := config.MustLoadConfig(new(config.Configuration), "config/config.toml")
exchangeConf := config.MustLoadConfig(new(ExchangeConf), "config/exchange.toml")
// etcdClient, err := clientv3.New(conf.Etcd)
// if err != nil {
// panic(err)
// }
// exit.AddHook(func() { _ = etcdClient.Close() }, exit.WithOrderTail())
// get market grpc client
// dis := discovery.NewEtcdDiscovery(etcdClient)
// resolver, err := dis.Resolver()
// if err != nil {
// panic(err)
// }
// Consul 配置
// consul 配置
cc := api.DefaultConfig()
cc.Address = conf.Consul.Address // Consul 地址
cc.Address = conf.Consul.Address
client, err := api.NewClient(cc)
if err != nil {
panic(fmt.Errorf("consul client error: %v", err))
@ -113,7 +100,6 @@ func main() {
if err := dis.Registry(grpcServer, register); err != nil {
panic(err)
}
zlog.Info("service registered to consul")
// run grpc server
go func() {

5
cmd/market/main.go

@ -27,9 +27,9 @@ func main() {
conf := config.MustLoadConfig(new(config.Configuration), "config/config.toml")
marketConf := config.MustLoadConfig(new(MarketConf), "config/market.toml")
// Consul 配置
// consul 配置
cc := api.DefaultConfig()
cc.Address = conf.Consul.Address // Consul 地址
cc.Address = conf.Consul.Address
client, err := api.NewClient(cc)
if err != nil {
panic(fmt.Errorf("consul client error: %v", err))
@ -72,7 +72,6 @@ func main() {
if err := dis.Registry(grpcServer, register); err != nil {
panic(err)
}
zlog.Info("service registered to consul")
// run grpc server
go func() {

4
cmd/sig-admin/main.go

@ -18,9 +18,9 @@ func main() {
// load config
conf := config.MustLoadConfig(new(config.Configuration), "config/config.toml")
// Consul 配置
// consul 配置
cc := api.DefaultConfig()
cc.Address = conf.Consul.Address // Consul 地址
cc.Address = conf.Consul.Address
client, err := api.NewClient(cc)
if err != nil {
panic(fmt.Errorf("consul client error: %v", err))

5
docker-compose.yml

@ -89,10 +89,11 @@ services:
sig-vm:
container_name: sig-vm
image: victoriametrics/victoria-metrics:v1.116.0
network_mode: host
volumes:
- "./fs/victoria-metrics-data:/victoria-metrics-data"
ports:
- 8428:8428
# ports:
# - 8428:8428
command: -dedup.minScrapeInterval=1s -retentionPeriod=99y
sig-questdb:

274
internal/exchange/exchange_grpc_server.go

@ -294,7 +294,7 @@ const (
HistoryKlineTsKey string = "history-kline-ts:%s:%s:%s" // exchange:sig-instid:interval
)
type initialKlineTask struct {
type fetchKlineTask struct {
inst types.TradeInstance
interval types.Interval
afterTs int64
@ -302,182 +302,184 @@ type initialKlineTask struct {
times int32 // 重试次数
}
func (t initialKlineTask) logKey() string {
func (t fetchKlineTask) logKey() string {
return fmt.Sprintf("%s:%s:%s:%d:%d", t.inst.Exchange, t.inst.InstId, t.interval, t.beforeTs, t.afterTs)
}
// initialKline 初始化交易产品历史k线数据
func (svc *ExchangeGrpcServer) initialKlines(exchange *Exchange, insts []types.TradeInstance) {
concurrent := max(8, runtime.NumCPU()*2)
// 记录成功和失败的交易产品
var success, failed []types.TradeInstance
for _, inst := range insts {
err := svc.initTradeInstanceKlines(exchange, inst)
if err != nil {
failed = append(failed, inst)
} else {
success = append(success, inst)
}
}
zlog.Infof("%d insts initial finished, success %d, failed %d", len(insts), len(success), len(failed))
}
var (
SingleTaskMaxFailTimes int32 = 10
)
// initTradeInstanceKlines 初始化交易产品历史k线数据
func (svc *ExchangeGrpcServer) initTradeInstanceKlines(exchange *Exchange, inst types.TradeInstance) (err error) {
// 并发数
concurrent := max(8, runtime.NumCPU()*2)
// 任务 channel
taskCh := make(chan initialKlineTask, concurrent)
taskCh := make(chan fetchKlineTask, concurrent)
retryTaskCh := make(chan fetchKlineTask, concurrent)
// 发布任务数, 成功任务数, 失败任务次数
var pubTasks, subTasks, failTasks atomic.Int32
ctx, cancel := context.WithCancel(context.Background())
// 任务生成
go func() {
for _, inst := range insts {
// for interval, intervalAdder := range types.SupportedIntervals {
interval := types.Interval1d
intervalAdder := types.SupportedIntervals[interval]
// history 未补全前, history写 kvdb ts mark, 补全后 ws live 写 ts mark
tsKey := fmt.Sprintf(HistoryKlineTsKey, inst.Exchange, inst.InstId, interval)
beforeTs, err := svc.kvdb.GetI64(context.Background(), tsKey)
if err != nil {
zlog.Error(err)
panic(err)
defer func() {
close(taskCh)
taskCh = nil
zlog.Infof("trade instance initial kline %s(%s), pub %d fetch tasks", pubTasks.Load(), inst.InstId, inst.Exchange)
}()
// for interval, intervalAdder := range types.SupportedIntervals {
interval := types.Interval1h
intervalAdder := types.SupportedIntervals[interval]
// history 未补全前, history写 kvdb ts mark, 补全后 ws live 写 ts mark
tsKey := fmt.Sprintf(HistoryKlineTsKey, inst.Exchange, inst.InstId, interval)
beforeTs, ex := svc.kvdb.GetI64(context.Background(), tsKey)
if ex != nil {
err = ex
zlog.Error(err)
cancel()
return
}
if beforeTs == 0 {
beforeTs = intervalAdder(KlineBefore0, -1)
}
for {
afterTs := intervalAdder(beforeTs, 101)
task := fetchKlineTask{
inst: inst,
interval: interval,
afterTs: afterTs,
beforeTs: beforeTs,
times: 0,
}
if beforeTs == 0 {
beforeTs = intervalAdder(KlineBefore0, -1)
// 发布任务
select {
case taskCh <- task:
pubTasks.Add(1)
case <-ctx.Done():
return
}
for {
afterTs := intervalAdder(beforeTs, 101)
taskCh <- initialKlineTask{
inst: inst,
interval: interval,
afterTs: afterTs,
beforeTs: beforeTs,
times: 0,
}
beforeTs = intervalAdder(afterTs, -1)
// 对比 ws 获取的实时k线
inst, ok := exchange.Insts.Load(inst.ExchangeInstId)
if !ok {
break
}
// 订阅完成
if inst.LiveMarkTs != 0 && beforeTs > inst.LiveMarkTs {
beforeTs = intervalAdder(afterTs, -1)
// 对比 ws 获取的实时k线
inst, ok := exchange.Insts.Load(inst.ExchangeInstId)
if !ok {
break
}
// 订阅完成
if inst.LiveMarkTs != 0 && beforeTs > inst.LiveMarkTs {
// history status -> ok
break
}
if beforeTs > time.Now().UnixMilli() {
if inst.LiveMarkTs != 0 {
// history status -> ok
break
}
if beforeTs > time.Now().UnixMilli() {
if inst.LiveMarkTs != 0 {
// history status -> ok
} else {
// live status -> not ok
}
break
} else {
// live status -> not ok
}
break
}
// }
}
close(taskCh)
}()
loc, _ := time.LoadLocation("Asia/Shanghai")
// 任务消费器 8协程并行
wg := new(sync.WaitGroup)
for range concurrent {
wg.Add(1)
go func() {
defer wg.Done()
var task fetchKlineTask
var ok bool
for {
task, ok := <-taskCh
select {
case <-ctx.Done():
return
case task, ok = <-taskCh:
case task, ok = <-retryTaskCh:
}
if !ok {
break
continue
}
if task.times > 0 {
zlog.Infof("retry fetch history kline task %d times: task -> %s", task.times, task.logKey())
}
interval, afterTs, beforeTs := task.interval, task.afterTs, task.beforeTs
klines, err := exchange.Fetcher.FetchHistoryKlines(context.Background(), task.inst.ExchangeInstId, interval, afterTs, beforeTs)
if err != nil {
zlog.Errorf("fetch history kline task error: task -> %s, err -> %v", task.logKey(), err)
// retry task
task.times++
taskCh <- task
return
}
if len(klines) == 0 {
continue
}
sts, ets := klines[0].Ts, klines[len(klines)-1].Ts
ss := time.UnixMilli(sts).In(loc).Format(times.FORMAT_DATE)
ee := time.UnixMilli(ets).In(loc).Format(times.FORMAT_DATE)
zlog.Infof("fetch interval %s %d~%d klines: ret=%d~%d, %d klines, %s~%s", interval, afterTs, beforeTs, sts, ets, len(klines), ss, ee)
// store to tsdb
err = svc.exchangeDataService.SaveKlines(task.inst, klines)
if err != nil {
zlog.Errorf("save history klines to tsdb error: task -> %s, err -> %v", task.logKey(), err)
if ex := svc.fetchTaskKlines(exchange, task); ex != nil {
failTasks.Add(1)
if task.times >= SingleTaskMaxFailTimes {
err = fmt.Errorf("task failed to many times %d, key: %s, err: %v", task.times, task.logKey(), ex)
cancel()
return
}
// retry task
task.times++
taskCh <- task
return
select {
case retryTaskCh <- task:
case <-ctx.Done():
return
}
} else {
// 任务都已执行成功结束
if subTasks.Add(1) >= pubTasks.Load() {
zlog.Infof("trade instance initial kline tasks success finished, %s(%s), pub %d, sub %d, fail %d", inst.InstId, inst.Exchange, pubTasks.Load(), subTasks.Load(), failTasks.Load())
cancel()
return
}
}
}
}()
}
wg.Wait()
return
}
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))
func (svc *ExchangeGrpcServer) fetchTaskKlines(exchange *Exchange, task fetchKlineTask) (err error) {
interval, afterTs, beforeTs := task.interval, task.afterTs, task.beforeTs
klines, err := exchange.Fetcher.FetchHistoryKlines(context.Background(), task.inst.ExchangeInstId, interval, afterTs, beforeTs)
if err != nil {
zlog.Errorf("fetch history kline task error: task -> %s, err -> %v", task.logKey(), err)
return
}
if len(klines) == 0 {
return
}
// // store to tsdb
// err = svc.exchangeDataService.SaveKlines(inst, klines)
// if err != nil {
// return
// }
loc, _ := time.LoadLocation("Asia/Shanghai")
// // set kvdb inst ts mark
// beforeTs = klines[0].Ts
// }
// // }
sts, ets := klines[0].Ts, klines[len(klines)-1].Ts
ss := time.UnixMilli(sts).In(loc).Format(times.FORMAT_DATE)
ee := time.UnixMilli(ets).In(loc).Format(times.FORMAT_DATE)
zlog.Infof("fetch interval %s %d~%d klines: ret=%d~%d, %d klines, %s~%s", interval, afterTs, beforeTs, sts, ets, len(klines), ss, ee)
// zlog.Infof("%s %s initial finished", inst.Exchange, inst.InstId)
// store to tsdb
err = svc.exchangeDataService.SaveKlines(task.inst, klines)
if err != nil {
zlog.Errorf("save history klines to tsdb error: task -> %s, err -> %v", task.logKey(), err)
return
}
return
}

3
pkg/grpc/discovery/consul_naming.go

@ -63,9 +63,10 @@ func (r *ConsulDiscovery) Registry(grpcServer grpc.ServiceRegistrar, register Se
},
}
if err = r.client.Agent().ServiceRegister(reg); err != nil {
zlog.Error("registry service error:", err)
zlog.Errorf("registry service %s error: %v", register.Name, err)
return
}
zlog.Infof("service %s registered to consul", register.Name)
return
}

Loading…
Cancel
Save