From 1c3ea73677d2705782c65dbb7be45b9faa647418 Mon Sep 17 00:00:00 2001
From: panlei <2799247126@qq.com>
Date: 星期四, 05 十二月 2019 18:33:02 +0800
Subject: [PATCH] ants协程池
---
main.go | 108 ++++++++++++++++++++++++++---------------------------
1 files changed, 53 insertions(+), 55 deletions(-)
diff --git a/main.go b/main.go
index 8c4179e..935a091 100644
--- a/main.go
+++ b/main.go
@@ -1,24 +1,24 @@
package main
import (
- "basic.com/dbapi.git"
- "basic.com/pubsub/protomsg.git"
- "basic.com/valib/deliver.git"
+
+ "flag"
+ "sync"
"net/http"
_ "net/http/pprof"
"plugin"
+
+ //"github.com/spf13/viper"
+ logger "github.com/alecthomas/log4go"
+ "github.com/panjf2000/ants/v2"
+
+ "basic.com/pubsub/protomsg.git"
+ "basic.com/valib/deliver.git"
"ruleprocess/insertdata"
"ruleprocess/labelFilter"
"ruleprocess/structure"
- "time"
-
- "basic.com/valib/logger.git"
- "flag"
- "fmt"
- "github.com/spf13/viper"
"ruleprocess/cache"
"ruleprocess/ruleserver"
- "sync"
)
var dbIp = flag.String("dbIp", "127.0.0.1", "dbserver ip")
@@ -33,18 +33,22 @@
// 鏃ュ織鍒濆鍖�
insertdata.Init(*env)
- var logFile = "./logger/"
- if viper.GetString("LogBasePath") != "" {
- logFile = viper.GetString("LogBasePath")
- }
- logFile = logFile + "ruleprocess.log"
- fmt.Println("鏃ュ織鍦板潃锛�",logFile)
- logger.Config(logFile, logger.DebugLevel)
- logger.SetSaveDays(7)
+ //var logFile = "./logger/"
+ //if viper.GetString("LogBasePath") != "" {
+ // logFile = viper.GetString("LogBasePath")
+ //}
+ //logFile = logFile + "ruleprocess.log"
+ //fmt.Println("鏃ュ織鍦板潃锛�",logFile)
+ //logger.Config(logFile, logger.DebugLevel)
+ //logger.SetSaveDays(7)
+ // log4go
+ logger.LoadConfiguration("./logger/log.xml")
logger.Info("鏃ュ織鍒濆鍖栨垚鍔燂紒")
+
}
func main() {
//fmt.Println("缂撳瓨鍒濆鍖栧畬鎴�",<- initchan)//dbserver鍒濆鍖栧畬姣�
+ defer ants.Release()
go func() {
http.ListenAndServe("0.0.0.0:8899",nil)
}()
@@ -52,15 +56,17 @@
wg := sync.WaitGroup{}
wg.Add(3)
- dbapi.Init(*dbIp, *dbPort)
go cache.Init(initchan, *dbIp, *surveyPort, *pubPort)
logger.Info("cache init completed!!!", <-initchan) //dbserver鍒濆鍖栧畬姣�
ruleserver.Init()
labelFilter.Init()
+
go ruleserver.TimeTicker()
go ruleserver.StartServer()
+
nReciever("ipc:///tmp/sdk-2-rules-process.ipc", deliver.PushPull, 1)
wg.Wait()
+
}
func nReciever(url string, m deliver.Mode, count int) {
c := deliver.NewServer(m, url)
@@ -68,50 +74,41 @@
}
func nRecvImpl(c deliver.Deliver, index int) {
-
var msg []byte
+ var wg1 sync.WaitGroup
+ p,_ := ants.NewPool(100)
+ syncCalculateSum := func() {
+ Task(msg)
+ wg1.Done()
+ }
+ wg1.Wait()
var err error
- //msgChan := make(chan []byte,100)
for {
- select {
- // case <-ctx.Done():
- // return
- default:
- msg, err = c.Recv()
- //msgChan <- msg
- if err != nil {
- logger.Info("recv error : ", err)
- fmt.Println("recv error : ", err)
- continue
- } else {
- //runtime.GOMAXPROCS(runtime.NumCPU())
- //logger.Debug("浣跨敤鐨刢pu涓暟锛�",runtime.NumCPU())
- //go func(msg []byte) {
- logger.Debug("褰撳墠鏃堕棿鎴筹細", time.Now().Unix())
- arg := structure.SdkDatas{}
- //paramFormat(msg, &arg)
- start := time.Now()
- m := CallParamFormat(msg, &arg)
- // 杩涜瑙勫垯澶勭悊鍒ゆ柇(鎵撲笂瑙勫垯鐨勬爣绛�)
- ruleserver.Judge(&arg, &m) // 鎶妔dkMessage浼犺繘鍘伙紝鏂逛究缂撳瓨鏁版嵁鏃舵嫾鍑轰竴涓猺esultMag
- // 鎶奱rg閲岀殑鎵撶殑鏍囩鎷垮嚭鏉ョ粰m鍐嶅皝瑁呬竴灞�
- resultMsg := structure.ResultMsg{SdkMessage: &m, RuleResult: arg.RuleResult}
- ruleserver.GetAttachInfo(resultMsg.SdkMessage)
- ruleEnd := time.Since(start)
- logger.Debug("瑙勫垯鍒ゆ柇瀹屾墍鐢ㄦ椂闂达細", ruleEnd)
- // 灏嗘墦瀹屾爣绛剧殑鏁版嵁鎻掑叆鍒癊S
- insertdata.InsertToEs(resultMsg)
- esEnd := time.Since(start)
- logger.Debug("鎻掑叆瀹孍s鎵�鐢ㄦ椂闂达細", esEnd)
- //浜嬩欢鎺ㄩ��
- labelFilter.PushSomthing(resultMsg)
- //}(msg)
- }
+ msg, err = c.Recv()
+ if err == nil {
+ wg1.Add(1)
+ _ = p.Submit(syncCalculateSum)
+ //go Task(msg)
}
}
}
+func Task (msg []byte) {
+ arg := structure.SdkDatas{}
+ //start := time.Now()
+ m := CallParamFormat(msg, &arg)
+ // 杩涜瑙勫垯澶勭悊鍒ゆ柇(鎵撲笂瑙勫垯鐨勬爣绛�)
+ ruleserver.Judge(&arg, &m) // 鎶妔dkMessage浼犺繘鍘伙紝鏂逛究缂撳瓨鏁版嵁鏃舵嫾鍑轰竴涓猺esultMag
+ // 鎶奱rg閲岀殑鎵撶殑鏍囩鎷垮嚭鏉ョ粰m鍐嶅皝瑁呬竴灞�
+ resultMsg := structure.ResultMsg{SdkMessage: &m, RuleResult: arg.RuleResult}
+ ruleserver.GetAttachInfo(resultMsg.SdkMessage)
+ // 灏嗘墦瀹屾爣绛剧殑鏁版嵁鎻掑叆鍒癊S
+ insertdata.InsertToEs(resultMsg)
+ //浜嬩欢鎺ㄩ��
+ labelFilter.PushSomthing(resultMsg)
+}
func CallParamFormat(msg []byte, args *structure.SdkDatas) protomsg.SdkMessage{
+ //logger.Info("鍛煎彨涓棿浠舵牸寮忓寲鏁版嵁")
p,err := plugin.Open("./algorithm/middleware.so")
if err != nil {
panic(err)
@@ -123,3 +120,4 @@
mess := f.(func(msg []byte, args *structure.SdkDatas)(protomsg.SdkMessage))(msg,args)
return mess
}
+
--
Gitblit v1.8.0