Senden von Nachrichten an und Empfangen von Nachrichten aus Azure Service Bus-Warteschlangen (Go)

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

Erstellen der Beispiel-App

Erstellen Sie zunächst ein neues Go-Modul.

  1. Erstellen Sie ein neues Verzeichnis für das Modul mit dem Namen service-bus-go-how-to-use-queues.

  2. Initialisieren Sie im azservicebus Verzeichnis 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/azservicebus
    
  3. Erstelle 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: