最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

golang微服務框架Kratos實現(xiàn)消息隊列

 更新時間:2025年12月15日 10:32:00   作者:喵了幾個咪  
本文介紹了在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):

  1. 在data層,使用消息隊列,但是在這個層,你只能在那生產消息,而不好消費消息。
  2. 將消息隊列的客戶端實現(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ù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • 淺析go中的map數(shù)據結構字典

    淺析go中的map數(shù)據結構字典

    golang中的map是一種數(shù)據類型,將鍵與值綁定到一起,底層是用哈希表實現(xiàn)的,可以快速的通過鍵找到對應的值。這篇文章主要介紹了go中的數(shù)據結構字典-map,需要的朋友可以參考下
    2019-11-11
  • golang處理TIFF圖像的實現(xiàn)示例

    golang處理TIFF圖像的實現(xiàn)示例

    本文介紹了在Go語言中處理TIFF圖像,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2025-03-03
  • goland設置顏色和字體的操作

    goland設置顏色和字體的操作

    這篇文章主要介紹了goland設置顏色和字體的操作方式,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-05-05
  • Go語言開發(fā)k8s之Service操作解析

    Go語言開發(fā)k8s之Service操作解析

    這篇文章主要為大家介紹了Go語言開發(fā)k8s之Service操作解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-06-06
  • Go Context庫 使用基本示例

    Go Context庫 使用基本示例

    在Go的http包中,每個請求由獨立的goroutine處理,這些goroutine可能需要訪問請求特定的數(shù)據或啟動其他服務,Context在Go語言中提供了一種方式來傳遞請求域的數(shù)據、取消信號和截止時間,本文介紹Go Context庫 使用基本示例,感興趣的朋友跟隨小編一起看看吧
    2024-09-09
  • golang?JSON序列化和反序列化示例詳解

    golang?JSON序列化和反序列化示例詳解

    通過使用Go語言的encoding/json包,你可以輕松地處理JSON數(shù)據,無論是在客戶端應用、服務器端應用還是其他類型的Go程序中,這篇文章主要介紹了golang?JSON序列化和反序列化,需要的朋友可以參考下
    2024-04-04
  • golang日志包logger的用法詳解

    golang日志包logger的用法詳解

    這篇文章主要介紹了golang日志包logger的用法詳解,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-05-05
  • 使用go操作redis的有序集合(zset)

    使用go操作redis的有序集合(zset)

    這篇文章主要介紹了使用go操作redis的有序集合(zset),具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-12-12
  • 詳解如何在Go語言中生成隨機種子

    詳解如何在Go語言中生成隨機種子

    這篇文章主要為大家詳細介紹了如何在Go語言中生成隨機種子,文中的示例代碼講解詳細,具有一定的借鑒價值,有需要的小伙伴可以參考一下
    2024-04-04
  • Go基本數(shù)據類型的具體使用

    Go基本數(shù)據類型的具體使用

    本文主要介紹了Go的基本數(shù)據類型,包括布爾類型、整數(shù)類型、浮點數(shù)類型、復數(shù)類型、字符串類型,具有一定的參考價值,感興趣的可以了解一下
    2023-11-11

最新評論

石嘴山市| 和静县| 绵阳市| 武川县| 张家港市| 宁阳县| 石台县| 九江市| 建宁县| 仪征市| 武乡县| 株洲县| 阿尔山市| 新营市| 榆中县| 清远市| 永春县| 龙陵县| 辽源市| 麟游县| 准格尔旗| 云阳县| 徐汇区| 泗阳县| 黑河市| 仙居县| 建宁县| 唐山市| 龙陵县| 邵武市| 凤城市| 玉林市| 南昌市| 乌鲁木齐县| 贡山| 墨脱县| 吉安县| 张家川| 永顺县| 靖江市| 厦门市|