Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Send and receive ordered messages

Última atualização: Jun 27, 2026

O ApsaraMQ for RocketMQ entrega e consome mensagens ordenadas seguindo rigorosamente a ordem FIFO (first-in, first-out). O código de exemplo a seguir demonstra como enviar e consumir mensagens ordenadas com o SDK cliente HTTP para Go.

Tipos de ordenação

As mensagens ordenadas dividem-se em duas categorias:

Tipo

Comportamento

Quando usar

Ordenação global

Todas as mensagens de um tópico seguem uma única sequência FIFO.

Quando cada mensagem precisa ser processada exatamente na ordem de envio, como na execução sequencial de comandos.

Ordenação por partição

As mensagens são distribuídas entre partições com base na chave de sharding. Cada partição mantém sua própria ordem FIFO.

Recomendado quando apenas mensagens que compartilham um agrupamento lógico, como o mesmo ID de pedido ou ID de usuário, exigem ordenação estrita.

Uma chave de sharding identifica a qual partição uma mensagem pertence. Ela difere da chave da mensagem: a chave da mensagem é um identificador de negócios usado para rastreamento, enquanto a chave de sharding controla a ordenação.

Para mais informações, consulte Mensagens ordenadas.

Como funciona a ordenação

A ordenação depende da coordenação entre o produtor e o consumidor.

Lado do envio: Utilize um único produtor e uma única thread para enviar mensagens. O broker preserva a sequência exata em que recebe as mensagens desse produtor.

Lado do consumo: Chame ConsumeMessageOrderly para puxar mensagens partição por partição. Dentro de cada partição, o próximo lote é entregue somente após o reconhecimento de todas as mensagens do lote atual. Caso o broker não receba um reconhecimento antes do prazo limite de NextConsumeTime, ele reentrega a mensagem não reconhecida.

Pré-requisitos

Conclua as configurações abaixo antes de executar o código de exemplo:

Enviar mensagens ordenadas

Importante

O broker determina a ordem das mensagens pela sequência em que um único produtor ou thread as envia. Se múltiplos produtores ou threads enviarem simultaneamente, a ordem de chegada ao broker pode divergir da ordem de negócios pretendida. Para garantir a ordenação, envie a partir de um único produtor usando uma única thread.

O código a seguir envia oito mensagens com ordenação por partição. A chave de sharding i % 2 distribui as mensagens entre duas partições; assim, mensagens com índices pares vão para uma partição e índices ímpares para outra.

Substitua os espaços reservados antes de executar o código:

Espaço reservado

Descrição

Exemplo

${HTTP_ENDPOINT}

Endpoint HTTP obtido na página Instance Details no console do ApsaraMQ for RocketMQ

http://xxx.mqrest.cn-hangzhou.aliyuncs.com

${TOPIC}

Tópico criado no console

OrderedTestTopic

${INSTANCE_ID}

ID da instância. Se a instância possuir um namespace, especifique o ID. Caso contrário, defina como uma string vazia. Verifique a página Instance Details para obter informações sobre o namespace.

MQ_INST_xxx

package main

import (
    "fmt"
    "time"
    "strconv"
    "os"

    "github.com/aliyunmq/mq-http-go-sdk"
)

func main() {
    // HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
    endpoint := "${HTTP_ENDPOINT}"
    // Set the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
    // environment variables before running this code.
    accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
    secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
    // Topic created in the ApsaraMQ for RocketMQ console.
    topic := "${TOPIC}"
    // Instance ID. If the instance has a namespace, specify the ID.
    // If not, set this to an empty string.
    instanceId := "${INSTANCE_ID}"

    client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")

    mqProducer := client.GetProducer(instanceId, topic)
    // Send 8 messages across 2 partitions.
    for i := 0; i < 8; i++ {
        msg := mq_http_sdk.PublishMessageRequest{
            MessageBody: "hello mq!",
            MessageTag:  "",
            Properties:  map[string]string{},
        }
        msg.MessageKey = "MessageKey"
        msg.Properties["a"] = strconv.Itoa(i)
        // Sharding key determines the partition.
        // Messages with the same sharding key are delivered in FIFO order.
        msg.ShardingKey = strconv.Itoa(i % 2)
        ret, err := mqProducer.PublishMessage(msg)

        if err != nil {
            fmt.Println(err)
            return
        } else {
            fmt.Printf("Publish ---->\n\tMessageId:%s, BodyMD5:%s, \n", ret.MessageId, ret.MessageBodyMD5)
        }
        time.Sleep(time.Duration(100) * time.Millisecond)
    }
}

Consumir mensagens ordenadas

Este código consome mensagens ordenadas usando long polling. O consumidor puxa mensagens partição por partição e reconhece cada lote antes de solicitar o próximo.

Substitua os espaços reservados pelos valores reais. Além dos espaços reservados listados na seção Enviar mensagens ordenadas, especifique o seguinte:

Espaço reservado

Descrição

Exemplo

${GROUP_ID}

Grupo de consumidores criado no console do ApsaraMQ for RocketMQ

GID_OrderedTest

package main

import (
    "fmt"
    "github.com/gogap/errors"
    "strings"
    "time"
    "os"

    "github.com/aliyunmq/mq-http-go-sdk"
)

func main() {
    // HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
    endpoint := "${HTTP_ENDPOINT}"
    // Set the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
    // environment variables before running this code.
    accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
    secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
    // Topic created in the ApsaraMQ for RocketMQ console.
    topic := "${TOPIC}"
    // Instance ID. If the instance has a namespace, specify the ID.
    // If not, set this to an empty string.
    instanceId := "${INSTANCE_ID}"
    // Consumer group created in the ApsaraMQ for RocketMQ console.
    groupId := "${GROUP_ID}"

    client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")

    mqConsumer := 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:
                {
                    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"+
                            "\tShardingKey: %s\n",
                            v.MessageId, v.PublishTime, v.MessageTag, v.ConsumedTimes,
                            v.FirstConsumeTime, v.NextConsumeTime, v.MessageBody, v.Properties, v.ShardingKey)
                    }

                    // Acknowledge all messages in this batch.
                    // If the broker does not receive an ACK before NextConsumeTime,
                    // it redelivers the message.
                    ackerr := mqConsumer.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(time.Duration(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(time.Duration(3) * time.Second)
                    }
                    endChan <- 1
                }
            case <-time.After(35 * time.Second):
                {
                    fmt.Println("Timeout of consumer message ??")
                    endChan <- 1
                }
            }
        }()

        // The consumer pulls partitionally ordered messages partition by partition.
        // Within each partition, the next batch is delivered only after all messages
        // in the current batch are acknowledged.
        // Long polling: if no message is available, the request is held on the broker
        // for up to the specified polling duration. The network timeout is 35 seconds.
        mqConsumer.ConsumeMessageOrderly(respChan, errChan,
            3, // Max messages per batch (up to 16).
            3, // Long polling duration in seconds (up to 30).
        )
        <-endChan
    }
}

Principais comportamentos de consumo

Comportamento

Detalhe

Ordenação no nível da partição

O consumidor pode puxar de múltiplas partições simultaneamente, mas as mensagens dentro de cada partição são sempre entregues na ordem de envio.

Barreira de reconhecimento em lote

O próximo lote de uma partição não é entregue até que todas as mensagens do lote anterior sejam reconhecidas.

Reentrega automática

Se o broker não receber nenhum reconhecimento antes de NextConsumeTime, ele reentrega a mensagem. Cada reentrega atribui um novo receipt handle com um timestamp único.

Long polling

Quando não há mensagens disponíveis, o broker retém a requisição por até a duração de polling especificada (3 segundos neste exemplo) antes de retornar uma resposta vazia. O timeout de rede é de 35 segundos.

Próximos passos

  • Mensagens ordenadas: Tipos de ordenação, garantias e considerações de design

  • Envie e receba outros tipos de mensagens usando o SDK HTTP para Go

  • Prepare o ambiente: Configure o ambiente de desenvolvimento Go