As mensagens normais são o tipo padrão no ApsaraMQ for RocketMQ e não possuem semântica especial de entrega. Diferentemente das mensagens agendadas, atrasadas, ordenadas e transacionais, as mensagens normais não apresentam comportamento adicional de envio.
Este tópico fornece códigos de exemplo para enviar e receber mensagens normais com o SDK cliente HTTP para Go.
Pré-requisitos
Antes de começar, verifique se você tem:
O SDK cliente HTTP para Go instalado. Para mais informações, consulte Prepare the environment.
Uma instância, um tópico e um grupo de consumidores do ApsaraMQ for RocketMQ criados no console. Para mais detalhes, acesse Create resources.
Um par de AccessKey da sua conta Alibaba Cloud. Consulte Create an AccessKey pair para obter instruções.
Parâmetros de configuração
Ambos os exemplos exigem os parâmetros listados abaixo. Substitua os placeholders pelos valores reais antes de executar o código.
|
Placeholder |
Descrição |
Onde encontrar |
|
|
O endpoint HTTP da sua instância |
Página Instance Details > seção HTTP Endpoint no console do ApsaraMQ for RocketMQ |
|
|
O tópico de destino das mensagens |
Crie-o no console. Cada tópico aceita apenas um tipo de mensagem |
|
|
O ID da sua instância |
Página Instance Details. Se a instância não tiver namespace, defina como nulo ou string vazia |
|
|
O ID do seu grupo de consumidores (apenas para consumo) |
Crie-o previamente no console |
O sistema lê as credenciais AccessKey das seguintes variáveis de ambiente:
ALIBABA_CLOUD_ACCESS_KEY_IDALIBABA_CLOUD_ACCESS_KEY_SECRET
Enviar mensagens normais
O exemplo a seguir inicializa um produtor e envia quatro mensagens para um tópico. Cada mensagem contém um corpo, uma tag opcional, uma chave de mensagem e propriedades personalizadas.
package main
import (
"fmt"
"os"
"strconv"
"time"
"github.com/aliyunmq/mq-http-go-sdk"
)
func main() {
endpoint := "<your-http-endpoint>"
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
topic := "<your-topic>"
instanceId := "<your-instance-id>"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
producer := client.GetProducer(instanceId, topic)
for i := 0; i < 4; i++ {
msg := mq_http_sdk.PublishMessageRequest{
MessageBody: "hello mq!",
MessageTag: "", // Optional: filter tag
Properties: map[string]string{}, // Optional: custom properties
}
msg.MessageKey = "MessageKey"
msg.Properties["a"] = strconv.Itoa(i)
ret, err := producer.PublishMessage(msg)
if err != nil {
fmt.Println(err)
return
}
fmt.Printf("Publish ---->\n\tMessageId:%s, BodyMD5:%s, \n",
ret.MessageId, ret.MessageBodyMD5)
time.Sleep(100 * time.Millisecond)
}
}
Pontos importantes:
A chamada
PublishMessageé síncrona e retornaMessageIdeMessageBodyMD5em caso de sucesso.Defina
MessageTagpara rotear mensagens a consumidores específicos que assinaram com um filtro de tag correspondente.Use
MessageKeypara atribuir um identificador personalizado à mensagem.
Consumir mensagens normais
Este exemplo inicializa um consumidor e busca mensagens continuamente via long polling. Após processar cada lote, o sistema confirma as mensagens para evitar reentrega.
package main
import (
"fmt"
"os"
"strings"
"time"
"github.com/aliyunmq/mq-http-go-sdk"
"github.com/gogap/errors"
)
func main() {
endpoint := "<your-http-endpoint>"
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
topic := "<your-topic>"
instanceId := "<your-instance-id>"
groupId := "<your-group-id>"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
consumer := client.GetConsumer(instanceId, topic, groupId, "")
for {
endChan := make(chan int)
respChan := make(chan mq_http_sdk.ConsumeMessageResponse)
errChan := make(chan error)
go func() {
select {
case resp := <-respChan:
// Process the batch
var handles []string
fmt.Printf("Consume %d messages---->\n", len(resp.Messages))
for _, v := range resp.Messages {
handles = append(handles, v.ReceiptHandle)
fmt.Printf("\tMessageID: %s, PublishTime: %d, MessageTag: %s\n"+
"\tConsumedTimes: %d, FirstConsumeTime: %d, NextConsumeTime: %d\n"+
"\tBody: %s\n"+
"\tProps: %s\n",
v.MessageId, v.PublishTime, v.MessageTag, v.ConsumedTimes,
v.FirstConsumeTime, v.NextConsumeTime, v.MessageBody, v.Properties)
}
// Acknowledge processed messages
ackerr := consumer.AckMessage(handles)
if ackerr != nil {
fmt.Println(ackerr)
if errAckItems, ok := ackerr.(errors.ErrCode).Context()["Detail"].([]mq_http_sdk.ErrAckItem); ok {
for _, errAckItem := range errAckItems {
fmt.Printf("\tErrorHandle:%s, ErrorCode:%s, ErrorMsg:%s\n",
errAckItem.ErrorHandle, errAckItem.ErrorCode, errAckItem.ErrorMsg)
}
} else {
fmt.Println("ack err =", ackerr)
}
time.Sleep(3 * time.Second)
} else {
fmt.Printf("Ack ---->\n\t%s\n", handles)
}
endChan <- 1
case err := <-errChan:
if strings.Contains(err.(errors.ErrCode).Error(), "MessageNotExist") {
fmt.Println("\nNo new message, continue!")
} else {
fmt.Println(err)
time.Sleep(3 * time.Second)
}
endChan <- 1
case <-time.After(35 * time.Second):
fmt.Println("Timeout of consumer message ??")
endChan <- 1
}
}()
// Long polling: wait up to 3 seconds for new messages, fetch up to 3 per batch
consumer.ConsumeMessage(respChan, errChan,
3, // Max messages per batch (up to 16)
3, // Long polling wait time in seconds (up to 30)
)
<-endChan
}
}
Funcionamento do long polling
No modo long polling, se não houver mensagens disponíveis, o broker retém a requisição pelo tempo especificado (até 30 segundos). Assim que uma mensagem chega, o broker responde imediatamente. Essa abordagem reduz respostas vazias em comparação ao short polling. O timeout padrão de rede é de 35 segundos.
Confirmação e reentrega
Após processar uma mensagem, chame AckMessage com o receipt handle para confirmar o consumo. Se o broker não receber a confirmação antes do prazo NextConsumeTime, ele reenvia a mensagem. Cada reentrega gera um novo receipt handle.
Nota: Se o receipt handle expirar antes da confirmação, o processo falhará. A resposta de erro inclui ErrorHandle, ErrorCode e ErrorMsg para cada handle com falha.