Skip to content
Toggle navigation
P
Projects
G
Groups
S
Snippets
Help
孙龙
/
go-api-behavior
This project
Loading...
Sign in
Toggle navigation
Go to a project
Project
Repository
Issues
0
Merge Requests
0
Pipelines
Wiki
Snippets
Settings
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Commit
133cfdd8
authored
Feb 24, 2020
by
孙龙
Browse files
Options
_('Browse Files')
Download
Email Patches
Plain Diff
init
parent
d8222df0
Expand all
Show whitespace changes
Inline
Side-by-side
Showing
6 changed files
with
193 additions
and
73 deletions
config/db.toml.demo
go.mod
go.sum
service/behavior/apiMsgService/LogSink.go
service/behavior/apiMsgService/apiMsgService.go
service/behavior/main.go
config/db.toml.demo
View file @
133cfdd8
...
@@ -3,9 +3,9 @@
...
@@ -3,9 +3,9 @@
dns="liexin:liexin#zsyM@tcp(192.168.2.232:3306)/liexin?parseTime=true"
dns="liexin:liexin#zsyM@tcp(192.168.2.232:3306)/liexin?parseTime=true"
[rabbitmq_ichunt]
[rabbitmq_ichunt]
queue_name="
send_buyer_mail
"
queue_name="
ichunt_monitor_user_behavior
"
routing_key="
send_buyer_mail
"
routing_key="
ichunt_monitor_user_behavior
"
exchange="ichunt_
order_msg
"
exchange="ichunt_
monitor_behavior
"
type="direct"
type="direct"
dns="amqp://guest:guest@192.168.2.232:5672/"
dns="amqp://guest:guest@192.168.2.232:5672/"
go.mod
View file @
133cfdd8
...
@@ -8,4 +8,5 @@ require (
...
@@ -8,4 +8,5 @@ require (
github.com/ichunt2019/go-msgserver v1.0.5 // indirect
github.com/ichunt2019/go-msgserver v1.0.5 // indirect
github.com/ichunt2019/logger v1.0.5
github.com/ichunt2019/logger v1.0.5
github.com/jmoiron/sqlx v1.2.0
github.com/jmoiron/sqlx v1.2.0
go.mongodb.org/mongo-driver v1.3.0
)
)
go.sum
View file @
133cfdd8
This diff is collapsed.
Click to expand it.
service/behavior/apiMsgService/LogSink.go
0 → 100644
View file @
133cfdd8
package
apiMsgService
import
(
"fmt"
//"fmt"
"go.mongodb.org/mongo-driver/mongo"
"context"
"go.mongodb.org/mongo-driver/mongo/options"
"time"
)
type
LogBatch
struct
{
Logs
[]
interface
{}
// 多条日志
}
// mongodb存储日志
type
LogSink
struct
{
client
*
mongo
.
Client
logCollection
*
mongo
.
Collection
logChan
chan
*
DataParams
autoCommitChan
chan
*
LogBatch
}
var
(
// 单例
G_logSink
*
LogSink
)
// 批量写入日志
func
(
logSink
*
LogSink
)
saveLogs
(
batch
*
LogBatch
)
{
logSink
.
logCollection
.
InsertMany
(
context
.
TODO
(),
batch
.
Logs
)
}
// 日志存储协程
func
(
logSink
*
LogSink
)
writeLoop
()
{
//select{
//case log := <- logSink.logChan:
// fmt.Println("从管道中读取日志")
// fmt.Println(log)
//}
var
(
log
*
DataParams
logBatch
*
LogBatch
// 当前的批次
commitTimer
*
time
.
Timer
//timeoutBatch *LogBatch // 超时批次
)
for
{
select
{
case
log
=
<-
logSink
.
logChan
:
fmt
.
Println
(
"新任务"
)
if
logBatch
==
nil
{
logBatch
=
&
LogBatch
{}
}
// 把新日志追加到批次中
logBatch
.
Logs
=
append
(
logBatch
.
Logs
,
log
)
// 如果批次满了, 就立即发送
if
len
(
logBatch
.
Logs
)
>=
10
{
// 发送日志
logSink
.
saveLogs
(
logBatch
)
// 清空logBatch
logBatch
=
nil
// 取消定时器
//commitTimer.Stop()
}
case
<-
time
.
NewTimer
(
1
*
time
.
Second
)
.
C
:
if
logBatch
==
nil
{
logBatch
=
&
LogBatch
{}
}
fmt
.
Println
(
"超时到期了"
)
fmt
.
Println
(
len
(
logBatch
.
Logs
))
logSink
.
saveLogs
(
logBatch
)
logBatch
=
nil
}
}
}
func
InitLogSink
()
(
err
error
)
{
var
(
client
*
mongo
.
Client
)
// 建立mongodb连接
clientOptions
:=
options
.
Client
()
.
ApplyURI
(
"mongodb://ichunt:huntmon6699@192.168.1.237:27017/ichunt?authMechanism=SCRAM-SHA-1"
)
if
client
,
err
=
mongo
.
Connect
(
context
.
TODO
(),
clientOptions
);
err
!=
nil
{
return
}
// 选择db和collection
G_logSink
=
&
LogSink
{
client
:
client
,
logCollection
:
client
.
Database
(
"ichunt"
)
.
Collection
(
"monitor_log"
),
logChan
:
make
(
chan
*
DataParams
,
1000
),
//autoCommitChan: make(chan *common.LogBatch, 1000),
}
// 启动一个mongodb处理协程
go
G_logSink
.
writeLoop
()
return
}
// 发送日志
func
(
logSink
*
LogSink
)
Append
(
jobLog
*
DataParams
)
{
select
{
case
logSink
.
logChan
<-
jobLog
:
default
:
// 队列满了就丢弃
}
}
\ No newline at end of file
service/behavior/apiMsgService/apiMsgService.go
View file @
133cfdd8
...
@@ -2,61 +2,72 @@ package apiMsgService
...
@@ -2,61 +2,72 @@ package apiMsgService
import
(
import
(
"encoding/json"
"encoding/json"
"fmt"
"go-api-behavior/util"
"time"
)
)
// 队列参数
func
ToJson
(
msg
string
)
(
err
error
){
type
DataParams
struct
{
err
=
json
.
Unmarshal
([]
byte
(
msg
),
&
util
.
MsgParams
)
InterfaceType
string
`json:"interface_type"`
return
err
AccessUrl
string
`json:"access_url"`
}
RequestParams
string
`json:"request_params"`
ErrMsg
string
`json:"err_msg"`
//写日志到管道
ErrCode
string
`json:"err_code"`
func
SaveElkLogChan
(
MsgParams
*
util
.
DataParams
){
Uid
string
`json:"uid"`
util
.
MsgChan
<-
MsgParams
UserName
string
`json:"user_name"`
UserIp
string
`json:"user_ip"`
Remakr
string
`json:"remark"`
CreateTime
int64
`json:"create_time"`
CreateTimeStr
string
`json:"create_time_str"`
}
}
func
SaveElkLog
(
LogReportDir
string
){
var
MsgParams
*
DataParams
var
(
logBatch
*
util
.
LogBatch
)
fileLog
,
_
:=
util
.
NewFileLogger
(
LogReportDir
)
fileLog
.
Init
()
for
{
select
{
case
util
.
MsgParams
=
<-
util
.
MsgChan
:
if
logBatch
==
nil
{
logBatch
=
&
util
.
LogBatch
{}
}
// 把新日志追加到批次中
logBatch
.
Logs
=
append
(
logBatch
.
Logs
,
util
.
MsgParams
)
// 如果批次满了, 就立即写入
if
len
(
logBatch
.
Logs
)
>
500
{
//写入日志
fmt
.
Println
(
"长度到了开始写入日志"
)
defer
fileLog
.
Close
()
fileLog
.
WriteLog
(
logBatch
)
// 清空logBatch
logBatch
=
nil
}
case
<-
time
.
NewTimer
(
10
*
time
.
Second
)
.
C
:
if
logBatch
==
nil
{
logBatch
=
&
util
.
LogBatch
{}
}
fmt
.
Println
(
"超时到期了"
)
fmt
.
Println
(
len
(
logBatch
.
Logs
))
//写入日志
fmt
.
Println
(
"超时了写入日志"
)
defer
fileLog
.
Close
()
fileLog
.
WriteLog
(
logBatch
)
// 清空logBatch
logBatch
=
nil
}
time
.
Sleep
(
time
.
Microsecond
*
10
)
//写日志到管道
}
func
SaveMsgToLogChan
(
msg
string
)
(
err
error
)
{
err
=
json
.
Unmarshal
([]
byte
(
msg
),
&
MsgParams
)
G_logSink
.
Append
(
MsgParams
)
return
err
}
}
//
//func SaveElkLog(LogReportDir string){
// var(
// logBatch *util.LogBatch
// )
// fileLog,_ := util.NewFileLogger(LogReportDir)
// fileLog.Init()
//
// for{
// select{
// case util.MsgParams = <- util.MsgChan:
// if logBatch == nil{
// logBatch = &util.LogBatch{}
// }
// // 把新日志追加到批次中
// logBatch.Logs = append(logBatch.Logs, util.MsgParams)
// // 如果批次满了, 就立即写入
// if len(logBatch.Logs) > 500{
// //写入日志
// fmt.Println("长度到了开始写入日志")
// defer fileLog.Close()
// fileLog.WriteLog(logBatch)
// // 清空logBatch
// logBatch = nil
// }
// case <- time.NewTimer(10*time.Second).C:
//
// if logBatch == nil{
// logBatch = &util.LogBatch{}
// }
// fmt.Println("超时到期了")
// fmt.Println(len(logBatch.Logs))
// //写入日志
// fmt.Println("超时了写入日志")
// defer fileLog.Close()
// fileLog.WriteLog(logBatch)
// // 清空logBatch
// logBatch = nil
// }
//
// time.Sleep(time.Microsecond*10)
// }
//}
\ No newline at end of file
service/behavior/main.go
View file @
133cfdd8
...
@@ -17,13 +17,12 @@ type RecvPro struct {
...
@@ -17,13 +17,12 @@ type RecvPro struct {
//// 实现消费者 消费消息失败 自动进入延时尝试 尝试3次之后入库db
//// 实现消费者 消费消息失败 自动进入延时尝试 尝试3次之后入库db
func
(
t
*
RecvPro
)
Consumer
(
dataByte
[]
byte
)
error
{
func
(
t
*
RecvPro
)
Consumer
(
dataByte
[]
byte
)
error
{
//
logger.Info
(string(dataByte))
//
fmt.Println
(string(dataByte))
err
:=
apiMsgService
.
ToJso
n
(
string
(
dataByte
))
err
:=
apiMsgService
.
SaveMsgToLogCha
n
(
string
(
dataByte
))
if
err
!=
nil
{
if
err
!=
nil
{
logger
.
Error
(
"%s"
,
err
)
logger
.
Error
(
"%s"
,
err
)
return
err
return
err
}
}
apiMsgService
.
SaveElkLogChan
(
util
.
MsgParams
)
////return errors.New("顶顶顶顶")
////return errors.New("顶顶顶顶")
return
nil
return
nil
}
}
...
@@ -37,17 +36,6 @@ func (t *RecvPro) FailAction(dataByte []byte) error {
...
@@ -37,17 +36,6 @@ func (t *RecvPro) FailAction(dataByte []byte) error {
}
}
//func initDb(dns string) (err error) {
// err = db.Init(dns)
// if err != nil {
// return
// }
//
// return
//}
var
ConfigDir
string
var
ConfigDir
string
var
LogDir
string
var
LogDir
string
var
LogReportDir
string
var
LogReportDir
string
...
@@ -70,7 +58,7 @@ func main() {
...
@@ -70,7 +58,7 @@ func main() {
//初始化配置文件
//初始化配置文件
util
.
Init
(
ConfigDir
)
util
.
Init
(
ConfigDir
)
util
.
MsgChan
=
make
(
chan
*
util
.
DataParams
,
10
)
//
util.MsgChan = make(chan *util.DataParams,10)
//
//
logConfig
:=
make
(
map
[
string
]
string
)
logConfig
:=
make
(
map
[
string
]
string
)
...
@@ -91,11 +79,11 @@ func main() {
...
@@ -91,11 +79,11 @@ func main() {
util
.
Configs
.
Rabbitmq_ichunt
.
Dns
,
util
.
Configs
.
Rabbitmq_ichunt
.
Dns
,
}
}
_
=
apiMsgService
.
InitLogSink
()
//go func(){
go
func
(){
// apiMsgService.SaveElkLog(LogReportDir)
apiMsgService
.
SaveElkLog
(
LogReportDir
)
//}()
}()
for
{
for
{
var
wg
sync
.
WaitGroup
var
wg
sync
.
WaitGroup
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment