Skip to content

Commit 63c95f8

Browse files
committed
迁移日志到MongoDB
1 parent a9f6967 commit 63c95f8

10 files changed

Lines changed: 222 additions & 77 deletions

File tree

backend/main.go

Lines changed: 9 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -34,39 +34,28 @@ func main() {
3434
panic(err)
3535
}
3636
log.Info("initialized config successfully")
37-
38-
// 初始化日志设置
39-
logLevel := viper.GetString("log.level")
40-
if logLevel != "" {
41-
log.SetLevelFromString(logLevel)
42-
}
43-
log.Info("initialized log config successfully")
44-
if viper.GetString("log.isDeletePeriodically") == "Y" {
45-
err := services.InitDeleteLogPeriodically()
46-
if err != nil {
47-
log.Error("init DeletePeriodically failed")
48-
panic(err)
49-
}
50-
log.Info("initialized periodically cleaning log successfully")
51-
} else {
52-
log.Info("periodically cleaning log is switched off")
53-
}
54-
5537
// 初始化Mongodb数据库
5638
if err := database.InitMongo(); err != nil {
5739
log.Error("init mongodb error:" + err.Error())
5840
debug.PrintStack()
5941
panic(err)
6042
}
61-
log.Info("initialized MongoDB successfully")
43+
log.Info("initialized mongodb successfully")
6244

6345
// 初始化Redis数据库
6446
if err := database.InitRedis(); err != nil {
6547
log.Error("init redis error:" + err.Error())
6648
debug.PrintStack()
6749
panic(err)
6850
}
69-
log.Info("initialized Redis successfully")
51+
log.Info("initialized redis successfully")
52+
53+
// 初始化日志设置
54+
if err := services.InitLogService(); err != nil {
55+
log.Error("init log error:" + err.Error())
56+
panic(err)
57+
}
58+
log.Info("initialized log successfully")
7059

7160
if model.IsMaster() {
7261
// 初始化定时任务

backend/model/log.go

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,23 @@
11
package model
22

33
import (
4+
"crawlab/database"
45
"crawlab/utils"
56
"github.com/apex/log"
7+
"github.com/globalsign/mgo/bson"
68
"os"
79
"runtime/debug"
10+
"time"
811
)
912

13+
type LogItem struct {
14+
Id bson.ObjectId `json:"_id" bson:"_id"`
15+
Message string `json:"msg" bson:"msg"`
16+
TaskId string `json:"task_id" bson:"task_id"`
17+
IsError bool `json:"is_error" bson:"is_error"`
18+
Ts time.Time `json:"ts" bson:"ts"`
19+
}
20+
1021
// 获取本地日志
1122
func GetLocalLog(logPath string) (fileBytes []byte, err error) {
1223

@@ -42,3 +53,27 @@ func GetLocalLog(logPath string) (fileBytes []byte, err error) {
4253
logBuf = logBuf[:n]
4354
return logBuf, nil
4455
}
56+
57+
func AddLogItem(l LogItem) error {
58+
s, c := database.GetCol("logs")
59+
defer s.Close()
60+
if err := c.Insert(l); err != nil {
61+
log.Errorf("insert log error: " + err.Error())
62+
debug.PrintStack()
63+
return err
64+
}
65+
return nil
66+
}
67+
68+
func GetLogItemList(filter interface{}, skip int, limit int, sortStr string) ([]LogItem, error) {
69+
s, c := database.GetCol("logs")
70+
defer s.Close()
71+
72+
var logItems []LogItem
73+
if err := c.Find(filter).Skip(skip).Limit(limit).Sort(sortStr).All(&logItems); err != nil {
74+
debug.PrintStack()
75+
return logItems, err
76+
}
77+
78+
return logItems, nil
79+
}

backend/model/task.go

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,19 @@ func (t *Task) GetResults(pageNum int, pageSize int) (results []interface{}, tot
109109
return
110110
}
111111

112+
func (t *Task) GetLogItems() (logItems []LogItem, err error) {
113+
query := bson.M{
114+
"task_id": t.Id,
115+
}
116+
117+
logItems, err = GetLogItemList(query, 0, constants.Infinite, "+_id")
118+
if err != nil {
119+
return logItems, err
120+
}
121+
122+
return logItems, nil
123+
}
124+
112125
func GetTaskList(filter interface{}, skip int, limit int, sortKey string) ([]Task, error) {
113126
s, c := database.GetCol("tasks")
114127
defer s.Close()

backend/routes/task.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -235,12 +235,12 @@ func DeleteTask(c *gin.Context) {
235235

236236
func GetTaskLog(c *gin.Context) {
237237
id := c.Param("id")
238-
logStr, err := services.GetTaskLog(id)
238+
logItems, err := services.GetTaskLog(id)
239239
if err != nil {
240240
HandleError(http.StatusInternalServerError, c, err)
241241
return
242242
}
243-
HandleSuccessData(c, logStr)
243+
HandleSuccessData(c, logItems)
244244
}
245245

246246
func GetTaskResults(c *gin.Context) {

backend/services/log.go

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"crawlab/utils"
1010
"encoding/json"
1111
"github.com/apex/log"
12+
"github.com/globalsign/mgo"
1213
"github.com/globalsign/mgo/bson"
1314
"github.com/spf13/viper"
1415
"io/ioutil"
@@ -162,5 +163,42 @@ func InitDeleteLogPeriodically() error {
162163

163164
c.Start()
164165
return nil
166+
}
167+
168+
func InitLogIndexes() error {
169+
s, c := database.GetCol("logs")
170+
defer s.Close()
171+
172+
_ = c.EnsureIndexKey("task_id")
173+
_ = c.EnsureIndex(mgo.Index{
174+
Key: []string{"$text:msg"},
175+
})
176+
177+
return nil
178+
}
179+
180+
func InitLogService() error {
181+
logLevel := viper.GetString("log.level")
182+
if logLevel != "" {
183+
log.SetLevelFromString(logLevel)
184+
}
185+
log.Info("initialized log config successfully")
186+
if viper.GetString("log.isDeletePeriodically") == "Y" {
187+
if err := InitDeleteLogPeriodically(); err != nil {
188+
log.Error("init DeletePeriodically failed")
189+
return err
190+
}
191+
log.Info("initialized periodically cleaning log successfully")
192+
} else {
193+
log.Info("periodically cleaning log is switched off")
194+
}
165195

196+
if model.IsMaster() {
197+
if err := InitLogIndexes(); err != nil {
198+
log.Errorf(err.Error())
199+
return err
200+
}
201+
}
202+
203+
return nil
166204
}

backend/services/task.go

Lines changed: 110 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package services
22

33
import (
4+
"bufio"
45
"crawlab/constants"
56
"crawlab/database"
67
"crawlab/entity"
@@ -160,15 +161,70 @@ func SetEnv(cmd *exec.Cmd, envs []model.Env, task model.Task, spider model.Spide
160161
return cmd
161162
}
162163

163-
func SetLogConfig(cmd *exec.Cmd, path string) error {
164-
fLog, err := os.Create(path)
164+
func SetLogConfig(cmd *exec.Cmd, t model.Task) error {
165+
//fLog, err := os.Create(path)
166+
//if err != nil {
167+
// log.Errorf("create task log file error: %s", path)
168+
// debug.PrintStack()
169+
// return err
170+
//}
171+
//cmd.Stdout = fLog
172+
//cmd.Stderr = fLog
173+
174+
// get stdout reader
175+
stdout, err := cmd.StdoutPipe()
176+
readerStdout := bufio.NewReader(stdout)
165177
if err != nil {
166-
log.Errorf("create task log file error: %s", path)
178+
log.Errorf("get stdout error: %s", err.Error())
167179
debug.PrintStack()
168180
return err
169181
}
170-
cmd.Stdout = fLog
171-
cmd.Stderr = fLog
182+
183+
// get stderr reader
184+
stderr, err := cmd.StderrPipe()
185+
readerStderr := bufio.NewReader(stderr)
186+
if err != nil {
187+
log.Errorf("get stdout error: %s", err.Error())
188+
debug.PrintStack()
189+
return err
190+
}
191+
192+
// read stdout
193+
go func() {
194+
for {
195+
line, err := readerStdout.ReadString('\n')
196+
if err != nil {
197+
break
198+
}
199+
line = strings.Replace(line, "\n", "", -1)
200+
_ = model.AddLogItem(model.LogItem{
201+
Id: bson.NewObjectId(),
202+
Message: line,
203+
TaskId: t.Id,
204+
IsError: false,
205+
Ts: time.Now(),
206+
})
207+
}
208+
}()
209+
210+
// read stderr
211+
go func() {
212+
for {
213+
line, err := readerStderr.ReadString('\n')
214+
line = strings.Replace(line, "\n", "", -1)
215+
if err != nil {
216+
break
217+
}
218+
_ = model.AddLogItem(model.LogItem{
219+
Id: bson.NewObjectId(),
220+
Message: line,
221+
TaskId: t.Id,
222+
IsError: true,
223+
Ts: time.Now(),
224+
})
225+
}
226+
}()
227+
172228
return nil
173229
}
174230

@@ -260,7 +316,7 @@ func ExecuteShellCmd(cmdStr string, cwd string, t model.Task, s model.Spider) (e
260316
cmd.Dir = cwd
261317

262318
// 日志配置
263-
if err := SetLogConfig(cmd, t.LogPath); err != nil {
319+
if err := SetLogConfig(cmd, t); err != nil {
264320
return err
265321
}
266322

@@ -566,54 +622,60 @@ func SpiderFileCheck(t model.Task, spider model.Spider) error {
566622
return nil
567623
}
568624

569-
func GetTaskLog(id string) (logStr string, err error) {
625+
func GetTaskLog(id string) (logItems []model.LogItem, err error) {
570626
task, err := model.GetTask(id)
571-
572627
if err != nil {
573628
return
574629
}
575630

576-
if IsMasterNode(task.NodeId.Hex()) {
577-
if !utils.Exists(task.LogPath) {
578-
fileDir, err := MakeLogDir(task)
579-
580-
if err != nil {
581-
log.Errorf(err.Error())
582-
}
583-
584-
fileP := GetLogFilePaths(fileDir, task)
585-
586-
// 获取日志文件路径
587-
fLog, err := os.Create(fileP)
588-
defer fLog.Close()
589-
if err != nil {
590-
log.Errorf("create task log file error: %s", fileP)
591-
debug.PrintStack()
592-
}
593-
task.LogPath = fileP
594-
if err := task.Save(); err != nil {
595-
log.Errorf(err.Error())
596-
debug.PrintStack()
597-
}
598-
599-
}
600-
// 若为主节点,获取本机日志
601-
logBytes, err := model.GetLocalLog(task.LogPath)
602-
if err != nil {
603-
log.Errorf(err.Error())
604-
logStr = err.Error()
605-
} else {
606-
logStr = utils.BytesToString(logBytes)
607-
}
608-
return logStr, err
609-
}
610-
// 若不为主节点,获取远端日志
611-
logStr, err = GetRemoteLog(task)
631+
logItems, err = task.GetLogItems()
612632
if err != nil {
613-
log.Errorf(err.Error())
614-
615-
}
616-
return logStr, err
633+
return logItems, err
634+
}
635+
636+
return logItems, nil
637+
638+
//if IsMasterNode(task.NodeId.Hex()) {
639+
// if !utils.Exists(task.LogPath) {
640+
// fileDir, err := MakeLogDir(task)
641+
//
642+
// if err != nil {
643+
// log.Errorf(err.Error())
644+
// }
645+
//
646+
// fileP := GetLogFilePaths(fileDir, task)
647+
//
648+
// // 获取日志文件路径
649+
// fLog, err := os.Create(fileP)
650+
// defer fLog.Close()
651+
// if err != nil {
652+
// log.Errorf("create task log file error: %s", fileP)
653+
// debug.PrintStack()
654+
// }
655+
// task.LogPath = fileP
656+
// if err := task.Save(); err != nil {
657+
// log.Errorf(err.Error())
658+
// debug.PrintStack()
659+
// }
660+
//
661+
// }
662+
// // 若为主节点,获取本机日志
663+
// logBytes, err := model.GetLocalLog(task.LogPath)
664+
// if err != nil {
665+
// log.Errorf(err.Error())
666+
// logStr = err.Error()
667+
// } else {
668+
// logStr = utils.BytesToString(logBytes)
669+
// }
670+
// return logStr, err
671+
//}
672+
//// 若不为主节点,获取远端日志
673+
//logStr, err = GetRemoteLog(task)
674+
//if err != nil {
675+
// log.Errorf(err.Error())
676+
//
677+
//}
678+
//return logStr, err
617679
}
618680

619681
func CancelTask(id string) (err error) {

devops/master/mongo-pv.yaml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ metadata:
88
spec:
99
storageClassName: manual
1010
capacity:
11-
storage: 10Gi
11+
storage: 3Gi
1212
accessModes:
1313
- ReadWriteOnce
1414
hostPath:
@@ -25,4 +25,4 @@ spec:
2525
- ReadWriteOnce
2626
resources:
2727
requests:
28-
storage: 10Gi
28+
storage: 3Gi

0 commit comments

Comments
 (0)