diff --git a/cmd/exchange/main.go b/cmd/exchange/main.go index 77a444c..b1311f7 100644 --- a/cmd/exchange/main.go +++ b/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() { diff --git a/cmd/market/main.go b/cmd/market/main.go index 2a849ef..029b40d 100644 --- a/cmd/market/main.go +++ b/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() { diff --git a/cmd/sig-admin/main.go b/cmd/sig-admin/main.go index 1c2966a..981ea08 100644 --- a/cmd/sig-admin/main.go +++ b/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)) diff --git a/docker-compose.yml b/docker-compose.yml index fe4ded4..a9b0fba 100644 --- a/docker-compose.yml +++ b/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: diff --git a/internal/exchange/exchange_grpc_server.go b/internal/exchange/exchange_grpc_server.go index 8b3a064..012e2ce 100644 --- a/internal/exchange/exchange_grpc_server.go +++ b/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 } diff --git a/pkg/grpc/discovery/consul_naming.go b/pkg/grpc/discovery/consul_naming.go index d122fcc..604fa50 100644 --- a/pkg/grpc/discovery/consul_naming.go +++ b/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 }