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