Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
In diesem Lernprogramm erfahren Sie, wie Sie Nachrichten mithilfe der Programmiersprache Go an Azure Service Bus-Warteschlangen senden und empfangen.
Azure Service Bus ist ein vollständig verwalteter Unternehmensnachrichtenbroker mit Nachrichtenwarteschlangen und Veröffentlichungs-/Abonnentenfunktionen. Service Bus wird verwendet, um Anwendungen und Dienste miteinander zu entkoppeln, was einen verteilten, zuverlässigen und leistungsstarken Nachrichtentransport ermöglicht.
Mit dem Azservicebus-Paket von Azure SDK für Go können Sie Nachrichten von Azure Service Bus senden und empfangen und die Programmiersprache Go verwenden.
Am Ende dieses Tutorials können Sie eine einzelne Nachricht oder einen Batch von Nachrichten an eine Warteschlange senden, Nachrichten empfangen und Nachrichten, die nicht verarbeitet werden, in die Warteschlange für unzustellbare Nachrichten verschieben.
Voraussetzungen
- Ein Azure-Abonnement. Sie können Ihre Visual Studio- oder MSDN-Abonnentenvorteile aktivieren oder sich für ein kostenloses Konto registrieren.
- Wenn Sie über keine Warteschlange verfügen, mit der Sie arbeiten können, folgen Sie den Schritten im Artikel Verwenden des Azure-Portals zum Erstellen einer Service Bus-Warteschlange, um eine Warteschlange zu erstellen.
- Go Version 1.18 oder höher
Erstellen der Beispiel-App
Erstellen Sie zunächst ein neues Go-Modul.
Erstellen Sie ein neues Verzeichnis für das Modul mit dem Namen
service-bus-go-how-to-use-queues.Initialisieren Sie im
azservicebusVerzeichnis das Modul, und installieren Sie die erforderlichen Pakete.go mod init service-bus-go-how-to-use-queues go get github.com/Azure/azure-sdk-for-go/sdk/azidentity go get github.com/Azure/azure-sdk-for-go/sdk/messaging/azservicebusErstelle eine neue Datei mit dem Namen
main.go.
Authentifizieren und Erstellen eines Clients
Erstellen Sie in der main.go Datei eine neue Funktion namens GetClient , und fügen Sie den folgenden Code hinzu:
func GetClient() *azservicebus.Client {
namespace, ok := os.LookupEnv("AZURE_SERVICEBUS_HOSTNAME") //ex: myservicebus.servicebus.windows.net
if !ok {
panic("AZURE_SERVICEBUS_HOSTNAME environment variable not found")
}
cred, err := azidentity.NewDefaultAzureCredential(nil)
if err != nil {
panic(err)
}
client, err := azservicebus.NewClient(namespace, cred, nil)
if err != nil {
panic(err)
}
return client
}
Die GetClient Funktion gibt ein neues azservicebus.Client Objekt zurück, das mithilfe eines Azure Service Bus-Namespaces und einer Anmeldeinformation erstellt wird. Der Namespace wird von der AZURE_SERVICEBUS_HOSTNAME Umgebungsvariable bereitgestellt. Und die Anmeldeinformationen werden mithilfe der azidentity.NewDefaultAzureCredential Funktion erstellt.
Für die lokale Entwicklung verwendet das DefaultAzureCredential Zugriffstoken von Azure CLI, das durch Ausführen des az login Befehls zur Authentifizierung bei Azure erstellt werden kann.
Tipp
Verwenden Sie die NewClientFromConnectionString-Funktion , um sich mit einer Verbindungszeichenfolge zu authentifizieren.
Senden von Nachrichten an eine Warteschlange
Erstellen Sie in der main.go Datei eine neue Funktion namens SendMessage , und fügen Sie den folgenden Code hinzu:
func SendMessage(message string, client *azservicebus.Client) {
sender, err := client.NewSender("myqueue", nil)
if err != nil {
panic(err)
}
defer sender.Close(context.TODO())
sbMessage := &azservicebus.Message{
Body: []byte(message),
}
err = sender.SendMessage(context.TODO(), sbMessage, nil)
if err != nil {
panic(err)
}
}
SendMessage verwendet zwei Parameter: eine Nachrichtenzeichenfolge und ein azservicebus.Client Objekt. Anschließend wird ein neues azservicebus.Sender Objekt erstellt und die Nachricht an die Warteschlange gesendet. Um Massennachrichten zu senden, fügen Sie die SendMessageBatch Funktion zu Ihrer main.go Datei hinzu.
func SendMessageBatch(messages []string, client *azservicebus.Client) {
sender, err := client.NewSender("myqueue", nil)
if err != nil {
panic(err)
}
defer sender.Close(context.TODO())
batch, err := sender.NewMessageBatch(context.TODO(), nil)
if err != nil {
panic(err)
}
for _, message := range messages {
if err := batch.AddMessage(&azservicebus.Message{Body: []byte(message)}, nil); err != nil {
panic(err)
}
}
if err := sender.SendMessageBatch(context.TODO(), batch, nil); err != nil {
panic(err)
}
}
SendMessageBatch verwendet zwei Parameter: ein Segment von Nachrichten und ein azservicebus.Client Objekt. Anschließend wird ein neues azservicebus.Sender Objekt erstellt und die Nachrichten an die Warteschlange gesendet.
Empfangen von Nachrichten aus einer Warteschlange
Nachdem Sie Nachrichten an die Warteschlange gesendet haben, können Sie sie mit dem azservicebus.Receiver Typ empfangen. Um Nachrichten aus einer Warteschlange zu empfangen, fügen Sie die GetMessage Funktion zu Ihrer main.go Datei hinzu.
func GetMessage(count int, client *azservicebus.Client) {
receiver, err := client.NewReceiverForQueue("myqueue", nil) //Change myqueue to env var
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), count, nil)
if err != nil {
panic(err)
}
for _, message := range messages {
body := message.Body
fmt.Printf("%s\n", string(body))
err = receiver.CompleteMessage(context.TODO(), message, nil)
if err != nil {
panic(err)
}
}
}
GetMessage verwendet ein azservicebus.Client Objekt und erstellt ein neues azservicebus.Receiver Objekt. Sie empfängt dann die Nachrichten aus der Warteschlange. Die Receiver.ReceiveMessages Funktion verwendet zwei Parameter: einen Kontext und die Anzahl der empfangenden Nachrichten. Die Receiver.ReceiveMessages Funktion gibt ein Objektsegment azservicebus.ReceivedMessage zurück.
Als Nächstes durchläuft eine for Schleife die Nachrichten und druckt den Nachrichtentext. Anschließend wird die CompleteMessage Funktion aufgerufen, um die Nachricht abzuschließen und aus der Warteschlange zu entfernen.
Nachrichten, die die Längenbeschränkungen überschreiten, an eine ungültige Warteschlange gesendet werden oder nicht erfolgreich verarbeitet werden, können an die Warteschlange für inaktive Briefe gesendet werden. Um Nachrichten an die Warteschlange für unzustellbare Nachrichten zu senden, fügen Sie die Funktion SendDeadLetterMessage zu Ihrer main.go-Datei hinzu.
func DeadLetterMessage(client *azservicebus.Client) {
deadLetterOptions := &azservicebus.DeadLetterOptions{
ErrorDescription: to.Ptr("exampleErrorDescription"),
Reason: to.Ptr("exampleReason"),
}
receiver, err := client.NewReceiverForQueue("myqueue", nil)
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), 1, nil)
if err != nil {
panic(err)
}
if len(messages) == 1 {
err := receiver.DeadLetterMessage(context.TODO(), messages[0], deadLetterOptions)
if err != nil {
panic(err)
}
}
}
DeadLetterMessage verwendet ein azservicebus.Client Objekt und ein azservicebus.ReceivedMessage Objekt. Anschließend sendet es die Nachricht an die Dead Letter Queue. Die Funktion verwendet zwei Parameter: einen Kontext und ein azservicebus.DeadLetterOptions Objekt. Die Receiver.DeadLetterMessage-Funktion gibt einen Fehler zurück, wenn die Nachricht nicht an die Dead-Letter-Queue gesendet werden kann.
Um Nachrichten aus der Warteschlange für tote Buchstaben zu empfangen, fügen Sie die ReceiveDeadLetterMessage Funktion zu Ihrer main.go Datei hinzu.
func GetDeadLetterMessage(client *azservicebus.Client) {
receiver, err := client.NewReceiverForQueue(
"myqueue",
&azservicebus.ReceiverOptions{
SubQueue: azservicebus.SubQueueDeadLetter,
},
)
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), 1, nil)
if err != nil {
panic(err)
}
for _, message := range messages {
fmt.Printf("DeadLetter Reason: %s\nDeadLetter Description: %s\n", *message.DeadLetterReason, *message.DeadLetterErrorDescription) //change to struct an unmarshal into it
err := receiver.CompleteMessage(context.TODO(), message, nil)
if err != nil {
panic(err)
}
}
}
GetDeadLetterMessage verwendet ein azservicebus.Client Objekt und erstellt ein neues azservicebus.Receiver Objekt mit Optionen für die Warteschlange für inaktive Buchstaben. Anschließend werden die Nachrichten aus der Warteschlange für unzustellbare Nachrichten empfangen. Die Funktion empfängt dann eine Nachricht aus der Warteschlange für unzustellbare Nachrichten. Anschließend werden der Grund und die Beschreibung für die Unzustellbarkeit dieser Nachricht gedruckt.
Beispielcode
package main
import (
"context"
"errors"
"fmt"
"os"
"github.com/Azure/azure-sdk-for-go/sdk/azcore/to"
"github.com/Azure/azure-sdk-for-go/sdk/azidentity"
"github.com/Azure/azure-sdk-for-go/sdk/messaging/azservicebus"
)
func GetClient() *azservicebus.Client {
namespace, ok := os.LookupEnv("AZURE_SERVICEBUS_HOSTNAME") //ex: myservicebus.servicebus.windows.net
if !ok {
panic("AZURE_SERVICEBUS_HOSTNAME environment variable not found")
}
cred, err := azidentity.NewDefaultAzureCredential(nil)
if err != nil {
panic(err)
}
client, err := azservicebus.NewClient(namespace, cred, nil)
if err != nil {
panic(err)
}
return client
}
func SendMessage(message string, client *azservicebus.Client) {
sender, err := client.NewSender("myqueue", nil)
if err != nil {
panic(err)
}
defer sender.Close(context.TODO())
sbMessage := &azservicebus.Message{
Body: []byte(message),
}
err = sender.SendMessage(context.TODO(), sbMessage, nil)
if err != nil {
panic(err)
}
}
func SendMessageBatch(messages []string, client *azservicebus.Client) {
sender, err := client.NewSender("myqueue", nil)
if err != nil {
panic(err)
}
defer sender.Close(context.TODO())
batch, err := sender.NewMessageBatch(context.TODO(), nil)
if err != nil {
panic(err)
}
for _, message := range messages {
err := batch.AddMessage(&azservicebus.Message{Body: []byte(message)}, nil)
if errors.Is(err, azservicebus.ErrMessageTooLarge) {
fmt.Printf("Message batch is full. We should send it and create a new one.\n")
}
}
if err := sender.SendMessageBatch(context.TODO(), batch, nil); err != nil {
panic(err)
}
}
func GetMessage(count int, client *azservicebus.Client) {
receiver, err := client.NewReceiverForQueue("myqueue", nil)
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), count, nil)
if err != nil {
panic(err)
}
for _, message := range messages {
body := message.Body
fmt.Printf("%s\n", string(body))
err = receiver.CompleteMessage(context.TODO(), message, nil)
if err != nil {
panic(err)
}
}
}
func DeadLetterMessage(client *azservicebus.Client) {
deadLetterOptions := &azservicebus.DeadLetterOptions{
ErrorDescription: to.Ptr("exampleErrorDescription"),
Reason: to.Ptr("exampleReason"),
}
receiver, err := client.NewReceiverForQueue("myqueue", nil)
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), 1, nil)
if err != nil {
panic(err)
}
if len(messages) == 1 {
err := receiver.DeadLetterMessage(context.TODO(), messages[0], deadLetterOptions)
if err != nil {
panic(err)
}
}
}
func GetDeadLetterMessage(client *azservicebus.Client) {
receiver, err := client.NewReceiverForQueue(
"myqueue",
&azservicebus.ReceiverOptions{
SubQueue: azservicebus.SubQueueDeadLetter,
},
)
if err != nil {
panic(err)
}
defer receiver.Close(context.TODO())
messages, err := receiver.ReceiveMessages(context.TODO(), 1, nil)
if err != nil {
panic(err)
}
for _, message := range messages {
fmt.Printf("DeadLetter Reason: %s\nDeadLetter Description: %s\n", *message.DeadLetterReason, *message.DeadLetterErrorDescription)
err := receiver.CompleteMessage(context.TODO(), message, nil)
if err != nil {
panic(err)
}
}
}
func main() {
client := GetClient()
fmt.Println("send a single message...")
SendMessage("firstMessage", client)
fmt.Println("send two messages as a batch...")
messages := [2]string{"secondMessage", "thirdMessage"}
SendMessageBatch(messages[:], client)
fmt.Println("\nget all three messages:")
GetMessage(3, client)
fmt.Println("\nsend a message to the Dead Letter Queue:")
SendMessage("Send message to Dead Letter", client)
DeadLetterMessage(client)
GetDeadLetterMessage(client)
}
Ausführen des Codes
Erstellen Sie vor dem Ausführen des Codes eine Umgebungsvariable mit dem Namen AZURE_SERVICEBUS_HOSTNAME. Legen Sie den Wert der Umgebungsvariablen auf den ServiceBus-Namespace fest.
export AZURE_SERVICEBUS_HOSTNAME=<YourServiceBusHostName>
Führen Sie als Nächstes den folgenden go run Befehl aus, um die App auszuführen:
go run main.go
Nächste Schritte
Weitere Informationen finden Sie unter den folgenden Links: