golang微服務框架Kratos實現(xiàn)消息隊列
什么是消息隊列
MQ就是消息隊列,是Message Queue的縮寫。消息隊列是一種通信方式。消息的本質就是一種數(shù)據結構。因為MQ把項目中的消息集中式的處理和存儲,所以MQ主要有解耦,并發(fā),和削峰的功能。
為什么要使用消息隊列
1. 異步
通常的微服務實現(xiàn)里面,都是通過RPC進行微服務之間的相互調用,這是同步的。如果消息隊列的話,可以實現(xiàn)異步的調用。至于異步有啥好處呢,主要是為了削峰。
2. 削峰
同步的調用會帶來一個問題:瞬時流量??蛻舻恼{用同步接口節(jié)奏,你是無法把控的,流量將會是忽高忽低的,猛的來一波,搞不好系統(tǒng)就崩了潰了。
如果消息隊列的話,可以實現(xiàn)異步的調用,并且可以實現(xiàn)削峰,請求進來,我先放到消息隊列里面去,慢慢消化掉,不至于猛的來一下,把系統(tǒng)擊垮。
3. 解耦
通常的微服務實現(xiàn)里面,都是通過RPC進行微服務之間的相互調用,那么意味著,你要做到一件事情,你必須要知道做事情的對方是誰。在微服務的世界里面,如果設計得不好,那就是一團糟的相互調用網絡,看得你暈暈的,運維會瘋,后面接手的開發(fā)人員也得瘋。
應用了消息隊列,你就只需要跟消息隊列這個代理打交道,單線聯(lián)系,關系簡單。我們只需要生產消息,消費消息,至于是誰消費的,誰生產的,完全不用去管它。架構上,就會清爽多了。所以,要對服務進行解耦,消息隊列是一個很好的選擇。
Kratos與消息隊列
Kratos現(xiàn)在的版本(v2.2.1)中,還沒有對消息隊列的直接支持,但是要運用還是容易的。官方有一個空殼示例代碼BeerShop,可以看到,在data層,使用Kafka的痕跡。
對于在Kratos微服務框架里面應用消息隊列,我認為有兩種方式可以實現(xiàn):
- 在data層,使用消息隊列,但是在這個層,你只能在那生產消息,而不好消費消息。
- 將消息隊列的客戶端實現(xiàn)為微服務的一個Server,然后在微服務的Service中消費消息和生產消息。
第一個方式的應用面不廣,更多的時候,第二種方式的應用面會更廣一些,我選擇了第二種方式。但是,Kratos官方并沒有支持這一種方式。故而,我只能夠自己動手實現(xiàn)了,我從另外一個微服務Go-Micro里面提取了其Broker的實現(xiàn),并且將其實現(xiàn)為Kratos框架里面的一個Server。事實證明,這樣是可行的,并且很好使。
你可能會問,為什么我不直接使用Go-Micro呢?因為Go-Micro是一個很重的微服務框架,盡管它的功能很豐富,幾乎支持了大部分的微服務需求。但是對于一個應用來說,我并不需要使用所有的技術棧、中間件,我只需要部分的技術棧。所以,我寧愿做加法,也不愿意去做減法。對于服務端來說,可控、可用、可維護是最重要的。極簡,是一個很好的選擇。另外,我還要腹誹一點,我從Go-Micro提取出來的Broker在測試的過程中發(fā)現(xiàn),都有一些瑕疵。
我實現(xiàn)的代碼,我放到了github:https://github.com/tx7do/kratos-transport。
它所支持的協(xié)議和消息隊列有:
- Kafka
- RabbitMQ
- NATS
- Redis
- MQTT
- WebSocket
基本上是夠用了。
kratos-transport的應用
它主要分為了3個部分:
1. Codec 編解碼器
編解碼器現(xiàn)在使用的是Kratos的編解碼器。
2. Broker 消息隊列客戶端
可以直接拿來使用,我拿Kafka舉例:
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"os"
"os/signal"
"syscall"
"github.com/go-kratos/kratos/v2/encoding"
"github.com/tx7do/kratos-transport/broker"
"github.com/tx7do/kratos-transport/broker/kafka"
)
const (
testBrokers = "localhost:9092"
testTopic = "test_topic"
testGroupId = "a-group"
)
type Hygrothermograph struct {
Humidity float64 `json:"humidity"`
Temperature float64 `json:"temperature"`
}
func registerHygrothermographHandler() broker.Handler {
return func(ctx context.Context, event broker.Event) error {
var msg *Hygrothermograph = nil
switch t := event.Message().Body.(type) {
case []byte:
msg = &Hygrothermograph{}
if err := json.Unmarshal(t, msg); err != nil {
return err
}
case string:
msg = &Hygrothermograph{}
if err := json.Unmarshal([]byte(t), msg); err != nil {
return err
}
case *Hygrothermograph:
msg = t
default:
return fmt.Errorf("unsupported type: %T", t)
}
if err := handleHygrothermograph(ctx, event.Topic(), event.Message().Headers, msg); err != nil {
return err
}
return nil
}
}
func handleHygrothermograph(_ context.Context, topic string, headers broker.Headers, msg *Hygrothermograph) error {
log.Printf("Headers: %+v, Humidity: %.2f Temperature: %.2f\n", headers, msg.Humidity, msg.Temperature)
return nil
}
func main() {
ctx := context.Background()
interrupt := make(chan os.Signal, 1)
signal.Notify(interrupt, syscall.SIGHUP, syscall.SIGINT, syscall.SIGTERM, syscall.SIGQUIT)
b := kafka.NewBroker(
broker.OptionContext(ctx),
broker.Addrs(testBrokers),
broker.Codec(encoding.GetCodec("json")),
)
_, err := b.Subscribe(testTopic,
registerHygrothermographHandler(),
func() broker.Any {
return &Hygrothermograph{}
},
broker.SubscribeContext(ctx),
broker.Queue(testGroupId),
)
if err != nil {
fmt.Println(err)
}
<-interrupt
}
3. Server 封裝給Kratos的Server實現(xiàn)
還是拿Kafka舉例:
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/go-kratos/kratos/v2"
"github.com/go-kratos/kratos/v2/encoding"
"github.com/tx7do/kratos-transport/broker"
"github.com/tx7do/kratos-transport/transport/kafka"
)
const (
testBrokers = "localhost:9092"
testTopic = "test_topic"
testGroupId = "a-group"
)
type Hygrothermograph struct {
Humidity float64 `json:"humidity"`
Temperature float64 `json:"temperature"`
}
func registerHygrothermographHandler() broker.Handler {
return func(ctx context.Context, event broker.Event) error {
var msg *Hygrothermograph = nil
switch t := event.Message().Body.(type) {
case []byte:
msg = &Hygrothermograph{}
if err := json.Unmarshal(t, msg); err != nil {
return err
}
case string:
msg = &Hygrothermograph{}
if err := json.Unmarshal([]byte(t), msg); err != nil {
return err
}
case *Hygrothermograph:
msg = t
default:
return fmt.Errorf("unsupported type: %T", t)
}
if err := handleHygrothermograph(ctx, event.Topic(), event.Message().Headers, msg); err != nil {
return err
}
return nil
}
}
func handleHygrothermograph(_ context.Context, topic string, headers broker.Headers, msg *Hygrothermograph) error {
log.Printf("Humidity: %.2f Temperature: %.2f\n", msg.Humidity, msg.Temperature)
return nil
}
func main() {
ctx := context.Background()
kafkaSrv := kafka.NewServer(
kafka.WithAddress([]string{testBrokers}),
kafka.WithCodec(encoding.GetCodec("json")),
)
_ = kafkaSrv.RegisterSubscriber(ctx,
testTopic, testGroupId, false,
registerHygrothermographHandler(),
func() broker.Any {
return &Hygrothermograph{}
})
app := kratos.New(
kratos.Name("kafka"),
kratos.Server(
kafkaSrv,
),
)
if err := app.Run(); err != nil {
log.Println(err)
}
}
另外再看一個例子,是Websocket的,它的應用其實也是很廣的:
package main
import (
"errors"
"fmt"
"log"
"github.com/go-kratos/kratos/v2"
"github.com/go-kratos/kratos/v2/encoding"
"github.com/tx7do/kratos-transport/transport/websocket"
)
var testServer *websocket.Server
const (
MessageTypeChat = iota + 1
)
type ChatMessage struct {
Type int `json:"type"`
Message string `json:"message"`
}
func main() {
wsSrv := websocket.NewServer(
websocket.WithAddress(":8800"),
websocket.WithPath("/ws"),
websocket.WithConnectHandle(handleConnect),
websocket.WithCodec(encoding.GetCodec("json")),
)
testServer = wsSrv
wsSrv.RegisterMessageHandler(MessageTypeChat,
func(sessionId websocket.SessionID, payload websocket.MessagePayload) error {
switch t := payload.(type) {
case *ChatMessage:
return handleChatMessage(sessionId, t)
default:
return errors.New("invalid payload type")
}
},
func() websocket.Any { return &ChatMessage{} },
)
app := kratos.New(
kratos.Name("websocket"),
kratos.Server(
wsSrv,
),
)
if err := app.Run(); err != nil {
log.Println(err)
}
}
func handleConnect(sessionId websocket.SessionID, register bool) {
if register {
fmt.Printf("%s connected\n", sessionId)
} else {
fmt.Printf("%s disconnect\n", sessionId)
}
}
func handleChatMessage(sessionId websocket.SessionID, message *ChatMessage) error {
fmt.Printf("[%s] Payload: %v\n", sessionId, message)
testServer.Broadcast(MessageTypeChat, *message)
return nil
}
具體的應用實例
我寫了一些實例代碼,并且都已經提交到了Kratos的examples代碼倉庫中去了。
kratos-cqrs
這是一個簡單的CQRS的實現(xiàn),主要就是拿了Kafka來消費來自于傳感器的遙感數(shù)據,然后把數(shù)據存儲到數(shù)據庫中去。
需要注意的是,這個實例并不夠完整,我并沒有實現(xiàn)MQTT的消費,沒有實現(xiàn)前端頁面等等。只實現(xiàn)了對Kafka的消費。
kratos-realtimemap
這是一個完整的物聯(lián)網相關的例子,有前端,有后端,可以完整的跑起來看。
通過MQTT接收一個開放的公交遙測數(shù)據源,然后通過REST和Websocket向前端發(fā)送數(shù)據,在地圖上展現(xiàn)出來車輛的軌跡、車輛的位置、車輛的速度、開關門狀態(tài)等等。
kratos-chatroom
最簡單的Websocket聊天室,客戶端發(fā)送消息,服務端接收之后立即廣播給其他客戶端。
中間件代碼
到此這篇關于golang微服務框架Kratos實現(xiàn)消息隊列的文章就介紹到這了,更多相關golang Kratos消息隊列內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

