From 80b0cf2ac6362591efb7306e88fa69e051fb9d10 Mon Sep 17 00:00:00 2001 From: liuxiaolong <736321739@qq.com> Date: 星期一, 25 十一月 2019 14:03:20 +0800 Subject: [PATCH] add rabbitmq --- go.sum | 2 + go.mod | 1 server.go | 90 ++++++++++++++++++++++++++++++++++++++------- 3 files changed, 79 insertions(+), 14 deletions(-) diff --git a/go.mod b/go.mod index 73d728b..21eafe7 100644 --- a/go.mod +++ b/go.mod @@ -8,4 +8,5 @@ github.com/rifflock/lfshook v0.0.0-20180920164130-b9218ef580f5 github.com/sirupsen/logrus v1.4.2 github.com/spf13/viper v1.4.0 + github.com/streadway/amqp v0.0.0-20190827072141-edfb9018d271 // indirect ) diff --git a/go.sum b/go.sum index 7add632..b382244 100644 --- a/go.sum +++ b/go.sum @@ -94,6 +94,8 @@ github.com/spf13/pflag v1.0.3/go.mod h1:DYY7MBk1bdzusC3SYhjObp+wFpr4gzcvqqNjLnInEg4= github.com/spf13/viper v1.4.0 h1:yXHLWeravcrgGyFSyCgdYpXQ9dR9c/WED3pg1RhxqEU= github.com/spf13/viper v1.4.0/go.mod h1:PTJ7Z/lr49W6bUbkmS1V3by4uWynFiR9p7+dSq/yZzE= +github.com/streadway/amqp v0.0.0-20190827072141-edfb9018d271 h1:WhxRHzgeVGETMlmVfqhRn8RIeeNoPr2Czh33I4Zdccw= +github.com/streadway/amqp v0.0.0-20190827072141-edfb9018d271/go.mod h1:AZpEONHx3DKn8O/DFsRAY58/XVQiIPMTMB1SddzLXVw= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= diff --git a/server.go b/server.go index a3298d9..7abaf58 100644 --- a/server.go +++ b/server.go @@ -9,7 +9,8 @@ "andriodServer/esutil" "andriodServer/extend/config" - log "andriodServer/log" + "andriodServer/log" + "github.com/streadway/amqp" ) var addr = flag.String("addr", "0.0.0.0", "The address to listen to") @@ -23,6 +24,8 @@ var IsHub = flag.String("hub", "hub", "hub is personIsHub=1") var Size = flag.Int("size", 100, "size default is 100") var env = flag.String("env", "config", "env set") +var mqIp = flag.String("mqIp", "172.17.50.245", "default mq ip") +var mqPort = flag.Int("mqPort", 5672, "default mq port") func main() { flag.Parse() @@ -30,27 +33,86 @@ log.SetLogLevel(*Level) config.Init(*env) fmt.Println(*port) - src := *addr + ":" + strconv.Itoa(*port) - listener, err := net.Listen("tcp", src) + //src := *addr + ":" + strconv.Itoa(*port) + //listener, err := net.Listen("tcp", src) + //if err != nil { + // log.Log.Errorln(err) + // return + //} + //log.Log.Infof("Listening on %s.\n", src) + + //fmt.Println("starting server success.") + //defer listener.Close() + + //connArr:=make([]net.Conn,0) + + //for { + // conn, err := listener.Accept()// + // + // connArr = append(connArr,conn) + // if err != nil { + // log.Log.Infoln("some connecion error: ", err) + // } + // go handleConnection(conn,connArr) + //} + + + mqAddr := "amqp://guest:guest@" + *mqIp + ":" + strconv.Itoa(*mqPort)+"/" + conn, err := amqp.Dial(mqAddr) if err != nil { - log.Log.Errorln(err) + log.Log.Infof("Failed to connect to RabbitMQ,err:",err) return } - log.Log.Infof("Listening on %s.\n", src) + defer conn.Close() - fmt.Println("starting server success.") - defer listener.Close() + ch, err := conn.Channel() + if err !=nil { + log.Log.Infof("Failed to open a channel,err:",err) + return + } + defer ch.Close() - connArr:=make([]net.Conn,0) + q, err := ch.QueueDeclare( + "alarm2Android", // name + false, // durable + false, // delete when unused + false, // exclusive + false, // no-wait + nil, // arguments + ) + if err !=nil { + log.Log.Infof("Failed to declare a queue,err:",err) + return + } + handleAlarmData2MQ(q, ch) +} +func handleAlarmData2MQ(queue amqp.Queue,ch *amqp.Channel) { + tick := time.NewTicker(3 * time.Second) + lastTime := time.Now() for { - conn, err := listener.Accept()// - - connArr = append(connArr,conn) - if err != nil { - log.Log.Infoln("some connecion error: ", err) + select { + case <-tick.C: + curTime := time.Now() + alarmData := esutil.PostAction(*sec, *Eurl, *Picurl, *IsHub, *Size, lastTime, curTime) + if alarmData != nil { + err := ch.Publish( + "", + queue.Name, + false, + false, + amqp.Publishing{ + ContentType: "text/plain", + Body: alarmData, + }) + if err !=nil { + log.Log.Infof("send to mq err:",err) + } else { + log.Log.Infof("send to mq success,len(body):",len(alarmData)) + } + } + lastTime = curTime } - go handleConnection(conn,connArr) } } -- Gitblit v1.8.0