feat: adjust async code
This commit is contained in:
@@ -14,16 +14,12 @@ import (
|
|||||||
"message-nest/routers"
|
"message-nest/routers"
|
||||||
"message-nest/service/cron_service"
|
"message-nest/service/cron_service"
|
||||||
"message-nest/service/env_service"
|
"message-nest/service/env_service"
|
||||||
"message-nest/service/send_message_service"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"sync"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
//go:embed web/dist/*
|
//go:embed web/dist/*
|
||||||
f embed.FS
|
f embed.FS
|
||||||
|
|
||||||
wg sync.WaitGroup
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
@@ -36,7 +32,7 @@ func init() {
|
|||||||
cron_service.Setup()
|
cron_service.Setup()
|
||||||
}
|
}
|
||||||
|
|
||||||
func GinServerUp() {
|
func main() {
|
||||||
gin.SetMode(setting.ServerSetting.RunMode)
|
gin.SetMode(setting.ServerSetting.RunMode)
|
||||||
|
|
||||||
routersInit := routers.InitRouter(f)
|
routersInit := routers.InitRouter(f)
|
||||||
@@ -60,13 +56,3 @@ func GinServerUp() {
|
|||||||
logrus.Error("Server err: ", err)
|
logrus.Error("Server err: ", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func main() {
|
|
||||||
wg.Add(1)
|
|
||||||
|
|
||||||
go GinServerUp()
|
|
||||||
go send_message_service.MessageConsumer(&wg)
|
|
||||||
|
|
||||||
wg.Wait()
|
|
||||||
fmt.Println("Server exit...")
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -28,4 +28,4 @@ var LogsCleanDefaultValueMap = map[string]string{
|
|||||||
const AboutSectionName = "about"
|
const AboutSectionName = "about"
|
||||||
|
|
||||||
// 限制goroutine的最大数量
|
// 限制goroutine的最大数量
|
||||||
var MaxSemaphore = make(chan struct{}, 2048)
|
var MaxSendTaskSemaphoreChan = make(chan string, 2048)
|
||||||
|
|||||||
@@ -47,7 +47,7 @@ func DoSendMassage(c *gin.Context) {
|
|||||||
appG.CResponse(http.StatusOK, "发送成功!", nil)
|
appG.CResponse(http.StatusOK, "发送成功!", nil)
|
||||||
} else {
|
} else {
|
||||||
// 异步发送
|
// 异步发送
|
||||||
send_message_service.Buffer <- msgService
|
msgService.AsyncSend()
|
||||||
appG.CResponse(http.StatusOK, "提交成功!", nil)
|
appG.CResponse(http.StatusOK, "提交成功!", nil)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,43 +0,0 @@
|
|||||||
package send_message_service
|
|
||||||
|
|
||||||
import (
|
|
||||||
"github.com/sirupsen/logrus"
|
|
||||||
"message-nest/pkg/constant"
|
|
||||||
|
|
||||||
"sync"
|
|
||||||
)
|
|
||||||
|
|
||||||
var maxBufferSize = 10
|
|
||||||
var Buffer = make(chan SendMessageService, maxBufferSize)
|
|
||||||
|
|
||||||
func DoSendTask(task SendMessageService, wg *sync.WaitGroup) {
|
|
||||||
defer wg.Done()
|
|
||||||
defer func() {
|
|
||||||
if r := recover(); r != nil {
|
|
||||||
logrus.Error("DoSendTask: Recovered from panic:", r)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
constant.MaxSemaphore <- struct{}{}
|
|
||||||
defer func() {
|
|
||||||
<-constant.MaxSemaphore
|
|
||||||
}()
|
|
||||||
|
|
||||||
go task.Send()
|
|
||||||
}
|
|
||||||
|
|
||||||
func MessageConsumer(wg *sync.WaitGroup) {
|
|
||||||
defer wg.Done()
|
|
||||||
|
|
||||||
for {
|
|
||||||
task, ok := <-Buffer
|
|
||||||
if !ok {
|
|
||||||
logrus.Error("MessageConsumer: Channel closed. Exiting.")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
wg.Add(1)
|
|
||||||
go DoSendTask(task, wg)
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
|
||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"github.com/sirupsen/logrus"
|
"github.com/sirupsen/logrus"
|
||||||
"message-nest/models"
|
"message-nest/models"
|
||||||
|
"message-nest/pkg/constant"
|
||||||
"message-nest/service/send_task_service"
|
"message-nest/service/send_task_service"
|
||||||
"message-nest/service/send_way_service"
|
"message-nest/service/send_way_service"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -14,6 +15,10 @@ const (
|
|||||||
SendFail = 0
|
SendFail = 0
|
||||||
)
|
)
|
||||||
|
|
||||||
|
//var logger = logrus.WithFields(logrus.Fields{
|
||||||
|
// "prefix": "[Message]",
|
||||||
|
//})
|
||||||
|
|
||||||
func errStrIsSuccess(errStr string) int {
|
func errStrIsSuccess(errStr string) int {
|
||||||
if errStr == "" {
|
if errStr == "" {
|
||||||
return SendSuccess
|
return SendSuccess
|
||||||
@@ -41,6 +46,23 @@ func (sm *SendMessageService) LogsAndStatusMark(errStr string, status int) {
|
|||||||
logrus.Infof("%s, 状态:%d", strings.Trim(errStr, "\n"), status)
|
logrus.Infof("%s, 状态:%d", strings.Trim(errStr, "\n"), status)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// AsyncSend 异步发送一个消息任务的所有实例
|
||||||
|
func (sm *SendMessageService) AsyncSend() {
|
||||||
|
defer func() {
|
||||||
|
if r := recover(); r != nil {
|
||||||
|
logrus.Error("AsyncSend: Recovered from panic:", r)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
// 限制并发异步发送数量
|
||||||
|
constant.MaxSendTaskSemaphoreChan <- ""
|
||||||
|
defer func() {
|
||||||
|
<-constant.MaxSendTaskSemaphoreChan
|
||||||
|
}()
|
||||||
|
|
||||||
|
go sm.Send()
|
||||||
|
}
|
||||||
|
|
||||||
// Send 发送一个消息任务的所有实例
|
// Send 发送一个消息任务的所有实例
|
||||||
func (sm *SendMessageService) Send() string {
|
func (sm *SendMessageService) Send() string {
|
||||||
sm.Status = SendSuccess
|
sm.Status = SendSuccess
|
||||||
|
|||||||
Reference in New Issue
Block a user