This commit is contained in:
jhx
2023-08-18 11:55:54 +08:00
commit b4fb4e8ed2
9 changed files with 693 additions and 0 deletions
+80
View File
@@ -0,0 +1,80 @@
package job
import (
"encoding/json"
"io/ioutil"
"log"
"os"
"gitee.com/jhx1219/machine-go/utils/common/task"
)
type JobConf struct {
JobList []JobListConf `json:"job_list"`
}
type JobListConf struct {
Name string `json:"name"`
}
type JobInfo struct {
Task *task.Task
}
type Job struct {
JobLogPath string
JobMap map[string]*JobInfo
}
func NewJob() *Job {
job := &Job{}
job.JobLogPath = os.Getenv("JOB_LOG_PATH")
job.JobMap = make(map[string]*JobInfo)
confPath := os.Getenv("CONF_PATH")
data, err := ioutil.ReadFile(confPath + "/job.json")
if err != nil {
log.Println("读取job文件失败: " + err.Error())
return job
}
// 解析 JSON 数据
var jobConf JobConf
err = json.Unmarshal(data, &jobConf)
if err != nil {
log.Println("解析job文件失败: " + err.Error())
return job
}
for _, conf := range jobConf.JobList {
t := task.NewTask(job.JobLogPath + "/" + conf.Name)
job.JobMap[conf.Name] = &JobInfo{
Task: t,
}
}
return job
}
func (j *Job) JobTaskServer() {
for _, job := range j.JobMap {
jobCopy := job
jobCopy.Task.TaskServer()
}
}
func (j *Job) LogJobTask(msg string) {
for _, job := range j.JobMap {
job.Task.LogTask(msg)
}
}
func (j *Job) SetJobFunc(job string, f func(string, *task.Task)) {
if _, ok := j.JobMap[job]; !ok {
return
}
j.JobMap[job].Task.TaskFunc = f
}
+91
View File
@@ -0,0 +1,91 @@
package mqtt
import (
"fmt"
"log"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
type MqttClient struct {
Client mqtt.Client
SubscribeList []SubscribeConf
}
type MqttConf struct {
Host string
Port string
Username string
Password string
ClientId string
}
type SubscribeConf struct {
Topic string
Handler func(mqtt.Client, mqtt.Message)
}
func NewMqttClient() *MqttClient {
c := &MqttClient{}
c.Client = nil
c.SubscribeList = make([]SubscribeConf, 0)
return c
}
func (c *MqttClient) SubTopic(topic string, f func(mqtt.Client, mqtt.Message)) {
c.SubscribeList = append(c.SubscribeList, SubscribeConf{
Topic: topic,
Handler: f,
})
}
func (c *MqttClient) Connect(conf MqttConf) {
host := conf.Host
port := conf.Port
username := conf.Username
password := conf.Password
clientId := conf.ClientId
opts := mqtt.NewClientOptions()
opts.AddBroker(fmt.Sprintf("tcp://%s:%s", host, port))
opts.SetAutoReconnect(true) // 自动重连
opts.SetUsername(username)
opts.SetPassword(password)
opts.SetClientID(clientId)
opts.SetOnConnectHandler(func(client mqtt.Client) {
log.Println("mqtt 连接成功")
// 创建订阅
log.Println("加载订阅...")
c.CreateSubscribe()
log.Println("加载订阅完成")
})
opts.SetConnectionLostHandler(func(client mqtt.Client, err error) {
log.Println("mqtt 连接断开")
})
client := mqtt.NewClient(opts)
if token := client.Connect(); token.Wait() && token.Error() != nil {
panic("创建mqtt客户端错误,错误信息: " + token.Error().Error())
}
c.Client = client
}
func (c *MqttClient) Publish(topic string, payload []byte) {
token := c.Client.Publish(topic, 0, false, payload)
token.Wait()
}
func (c *MqttClient) CreateSubscribe() {
for _, conf := range c.SubscribeList {
if token := c.Client.Subscribe(conf.Topic, 1, conf.Handler); token.Wait() && token.Error() != nil {
log.Fatal(token.Error())
}
log.Printf("订阅 %s 完成\n", conf.Topic)
}
}
+39
View File
@@ -0,0 +1,39 @@
package receive
import (
"gitee.com/jhx1219/machine-go/utils/common/task"
)
type Receive struct {
Task *task.Task
}
func NewReceive(t *task.Task) *Receive {
t.SetTaskFunc(func(msg string, task *task.Task) {
//log.Println(msg)
})
return &Receive{
Task: t,
}
}
func (r *Receive) SetTaskFunc(f func(string, *task.Task)) {
r.Task.SetTaskFunc(f)
}
func (r *Receive) LogReceice(msg string) error {
err := r.Task.LogTask(msg)
if err != nil {
return err
}
return nil
}
func (r *Receive) ReceiveTaskServer() {
r.Task.TaskServer()
}
+70
View File
@@ -0,0 +1,70 @@
package request
import (
"bytes"
"errors"
"fmt"
"io/ioutil"
"log"
"net/http"
"time"
)
type Request struct {
Timeout time.Duration
MaxRetry int
}
func NewRequest() *Request {
return &Request{
Timeout: time.Second * 2,
}
}
func (r *Request) SetTimeout(t time.Duration) {
r.Timeout = t
}
func (r *Request) SetMaxRetry(retry int) {
r.MaxRetry = retry
}
func (r *Request) Post(api string, data []byte) (string, error) {
client := &http.Client{
Timeout: r.Timeout,
}
req, err := http.NewRequest("POST", api, bytes.NewBuffer(data))
if err != nil {
return "", errors.New("Error creating request:" + err.Error())
}
req.Header.Set("Content-Type", "application/json")
retry := 0
reqDo:
resp, err := client.Do(req)
if err != nil {
if r.MaxRetry > 0 {
if retry < r.MaxRetry {
retry = retry + 1
log.Printf("请求api第[%d], 最大重试次数[%d]", retry, r.MaxRetry)
goto reqDo
}
}
return "", errors.New("Error sending request:" + err.Error())
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
return "", fmt.Errorf("Error StatusCode: %d", resp.StatusCode)
}
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return "", errors.New("Error read body: " + err.Error())
}
return string(body), nil
}
+253
View File
@@ -0,0 +1,253 @@
package task
import (
"bufio"
"log"
"os"
"path/filepath"
"strconv"
"strings"
"sync"
"time"
)
type Task struct {
Path string
FileMutex *sync.Mutex
TaskFunc func(string, *Task)
RuningTaskPath map[string]int
}
func NewTask(path string) *Task {
var fileMutex sync.Mutex
return &Task{
Path: path,
FileMutex: &fileMutex,
TaskFunc: nil,
RuningTaskPath: make(map[string]int),
}
}
func (task *Task) SetTaskFunc(f func(string, *Task)) {
task.TaskFunc = f
}
func (task *Task) LogTask(msg string) error {
task.FileMutex.Lock()
defer task.FileMutex.Unlock()
path := task.Path
task.checkOrCreatePath(path)
currentTime := time.Now().Unix()
filePath := path + "/" + strconv.FormatInt(currentTime, 10) + ".log" // 目标文件路径
file, err := os.OpenFile(filePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
if err != nil {
log.Println("打开task文件失败:", err)
return err
}
defer file.Close()
// 写入内容到文件
_, err = file.WriteString(msg + "\n")
if err != nil {
log.Println("写入task文件失败:", err)
return err
}
return nil
}
func (task *Task) LogFailTask(msg string) error {
task.FileMutex.Lock()
defer task.FileMutex.Unlock()
path := task.Path + "/fail/"
task.checkOrCreatePath(path)
currentTime := time.Now()
currentHour := time.Date(currentTime.Year(), currentTime.Month(), currentTime.Day(), currentTime.Hour(), 0, 0, 0, currentTime.Location())
timestamp := currentHour.Unix()
filePath := path + strconv.FormatInt(timestamp, 10) + ".log" // 目标文件路径
file, err := os.OpenFile(filePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
if err != nil {
log.Println("打开task文件失败:", err)
return err
}
defer file.Close()
// 写入内容到文件
_, err = file.WriteString(msg + "\n")
if err != nil {
log.Println("写入task文件失败:", err)
return err
}
return nil
}
func (task *Task) checkOrCreatePath(path string) {
dir := path
// 检查任务目录是否存在
_, err := os.Stat(dir)
if os.IsNotExist(err) {
// 任务目录不存在,创建任务目录
err := os.MkdirAll(dir, 0755)
if err != nil {
log.Println(task.Path+" 目录创建失败:", err)
return
}
log.Println(task.Path + " 目录创建成功")
}
}
func (task *Task) TaskServer() {
// 扫描接收文件协程
go func() {
task.ScanTask()
}()
}
func (task *Task) ScanTask() {
path := task.Path
log.Println("扫描任务目录: " + path)
for {
currentTime := time.Now().Unix()
dirRead, err := os.Open(path)
if err != nil {
time.Sleep(1 * time.Second)
continue
}
for {
files, err := dirRead.Readdir(10) // 每次读取10个文件
if err != nil && len(files) == 0 {
// 读取完所有文件或发生错误
time.Sleep(500 * time.Microsecond)
break
}
for _, file := range files {
fileName := file.Name()
ext := filepath.Ext(fileName)
fileNameWithoutExt := strings.TrimSuffix(fileName, ext)
fileTime, err := strconv.ParseInt(fileNameWithoutExt, 10, 64)
if err != nil {
continue // 跳过无法转换为时间戳的文件名
}
if ext == ".log" {
if fileTime < currentTime-1 {
filePath := path + "/" + fileName
log.Println("扫描到文件: " + filePath)
newFilePath := path + "/" + fileNameWithoutExt + ".scanned"
err := os.Rename(filePath, newFilePath)
if err != nil {
log.Println(filePath + " 文件重命名失败: " + err.Error())
}
task.HandleTask(path + "/" + fileNameWithoutExt)
}
}
if ext == ".scanned" {
if fileTime < currentTime-600 {
filePath := path + "/" + fileName
// 非执行中的十分钟前未处理的重新处理
if _, exists := task.RuningTaskPath[filePath]; !exists {
log.Println("扫描到超时scanned需重新执行文件: " + filePath)
task.HandleTask(path + "/" + fileNameWithoutExt)
}
}
}
}
}
dirRead.Close()
}
}
func (task *Task) HandleTask(path string) {
filePath := path + ".scanned"
task.RuningTaskPath[path] = 1
log.Println("开始处理文件: " + filePath)
// 判断文件是否存在
_, err := os.Stat(filePath)
if os.IsNotExist(err) {
log.Println(filePath + " 文件不存在")
return
}
// 打开文件
file, err := os.Open(filePath)
if err != nil {
log.Println(filePath + " 文件打开失败: " + err.Error())
return
}
// 自定义缓冲区大小为 20 * 1MB
const bufferSize = 20 * 1024 * 1024
// 逐行读取文件内容
scanner := bufio.NewScanner(file)
// 设置缓冲区大小
buf := make([]byte, bufferSize)
scanner.Buffer(buf, bufferSize)
i := 0 // 消息数计数器
wg := sync.WaitGroup{}
lineLimitChan := make(chan bool, 10) // 限制协程数量
for scanner.Scan() {
lineLimitChan <- true
line := scanner.Text()
if task.TaskFunc != nil {
wg.Add(1)
go func() {
task.TaskFunc(line, task)
<-lineLimitChan
wg.Done()
}()
}
i = i + 1
}
wg.Wait()
if err := scanner.Err(); err != nil {
err := os.Rename(filePath, path+".error")
if err != nil {
log.Println(filePath + " 文件重命名失败: " + err.Error())
}
log.Println(filePath + " 文件扫描出错: " + err.Error())
return
}
// 文件扫描完毕
file.Close()
log.Printf(filePath+" 文件处理完毕, 共处理[%d]条信息", i)
err = os.Remove(filePath)
if err != nil {
log.Println(filePath + " 文件删除失败: " + err.Error())
return
}
log.Println(filePath + " 文件处理完成后删除")
delete(task.RuningTaskPath, path)
}