Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Métodos e parâmetros

Última atualização: Jun 27, 2026

O SDK Java do ApsaraMQ for RocketMQ envia e consome mensagens nos modos push ou pull. Este tópico descreve os métodos disponíveis e os parâmetros de configuração.

Modos push e pull

O ApsaraMQ for RocketMQ entrega mensagens aos consumidores em um dos dois modos:

  • Modo push: o broker envia mensagens aos consumidores assim que ficam disponíveis. Esse modo oferece suporte ao consumo em lote, entregando várias mensagens em um único lote.

  • Modo pull: os consumidores consultam o broker sob demanda para obter novas mensagens. Esse modo proporciona controle direto sobre quais partições ler, quando buscar mensagens e como gerenciar offsets.

Importante

O modo pull exige uma instância Enterprise Platinum Edition.

Parâmetros de conexão

Configure os seguintes parâmetros tanto para produtores quanto para consumidores.

Parâmetro

Descrição

Padrão

NAMESRV_ADDR

Endpoint TCP da instância do ApsaraMQ for RocketMQ. Encontre essa informação na página Instance Details no console.

--

AccessKey

AccessKey ID usado como identificador exclusivo para autenticação. Consulte Crie um par de AccessKey.

--

SecretKey

Segredo do AccessKey usado como senha para autenticação. Consulte Crie um par de AccessKey.

--

OnsChannel

Canal do usuário. Defina como CLOUD para usuários do CloudTmall.

ALIYUN

Métodos do produtor

A interface do produtor fornece os seguintes métodos:

messagesendinterface

Parâmetros do produtor

Parâmetro

Descrição

Padrão

Unidade

SendMsgTimeoutMillis

Tempo limite para envio de uma mensagem.

--

Milissegundos

CheckImmunityTimeInSeconds

Tempo mínimo de espera antes que o broker verifique pela primeira vez o status de uma mensagem transacional.

--

Segundos

shardingKey

Chave de partição que determina qual partição recebe uma mensagem ordenada.

--

--

Métodos do consumidor

Modo push

A interface do consumidor push fornece os seguintes métodos:

consumeinterface

Modo pull

A interface do consumidor pull fornece os seguintes métodos:

pull_consumer

Interface PullConsumer

public interface PullConsumer extends Admin {

    /**
     * Returns all partitions for the specified topic.
     * Call this method only after the pull consumer has started.
     */
    Set<TopicPartition> topicPartitions(String topic);

    /**
     * Manually assigns partitions to this consumer. Automatic rebalancing
     * is not performed -- make sure all partitions are covered across
     * your consumers. Calling this method again replaces the previous
     * assignment entirely.
     */
    void assign(Collection<TopicPartition> topicPartitions);

    /**
     * Polls for messages. Returns up to maxBatchMessageCount messages,
     * blocking for at most the specified timeout (in milliseconds).
     */
    List<Message> poll(long timeout);

    /**
     * Resets the consumer offset of a partition to a specific position.
     * The offset must be between the partition's minimum and maximum
     * offset. Call this method only after the consumer has started and
     * has been assigned the target partition.
     */
    void seek(TopicPartition topicPartition, long offset);

    /**
     * Resets the consumer offset of a partition to the earliest available
     * message. Requires a started consumer assigned to the target
     * partition.
     */
    void seekToBeginning(TopicPartition topicPartition);

    /**
     * Resets the consumer offset of a partition to the latest message.
     * Requires a started consumer assigned to the target partition.
     */
    void seekToEnd(TopicPartition topicPartition);

    /**
     * Pauses message consumption on the specified partitions.
     */
    void pause(Collection<TopicPartition> topicPartitions);

    /**
     * Resumes message consumption on previously paused partitions.
     */
    void resume(Collection<TopicPartition> topicPartitions);

    /**
     * Returns the offset of the first message stored at or after the
     * given timestamp in the specified partition. The timestamp refers
     * to when the broker stored the message, not when it was sent.
     */
    Long offsetForTimestamp(TopicPartition topicPartition, Long timestamp);

    /**
     * Returns the latest consumer offset for the specified partition.
     */
    Long committed(TopicPartition topicPartition);

    /**
     * Commits the current consumer offsets synchronously. Offsets are
     * first synced to the local client, then written to the broker
     * asynchronously.
     */
    void commitSync();

    /**
     * Listener interface for partition change events.
     */
    interface TopicPartitionChangeListener {
        /**
         * Called when the partitions of a topic change -- for example,
         * after broker scaling adjusts the partition count.
         */
        void onChanged(Set<TopicPartition> topicPartitions);
    }

    /**
     * Registers a listener that fires when the partition count for a
     * topic changes (for example, due to broker scaling). By default,
     * the maximum callback delay is 5 seconds.
     */
    void registerTopicPartitionChangedListener(String topic,
            TopicPartitionChangeListener callback);
}

Parâmetros do consumidor

Parâmetros comuns

Os parâmetros abaixo se aplicam tanto a consumidores push quanto pull.

Parâmetro

Descrição

Padrão

Valores válidos

Unidade

GROUP_ID

ID do grupo de consumidores criado no console do ApsaraMQ for RocketMQ. Consulte Termos.

--

--

--

MessageModel

Modelo de consumo. Valores válidos: CLUSTERING (consumo em cluster) e BROADCASTING (consumo por broadcast).

CLUSTERING

CLUSTERING, BROADCASTING

--

ConsumeThreadNums

Quantidade de threads que o consumidor utiliza para processar mensagens.

20

--

--

MaxReconsumeTimes

Número máximo de tentativas de reconsumo quando falha o consumo de uma mensagem.

16

--

--

ConsumeTimeout

Tempo máximo permitido para consumir uma única mensagem. Se ultrapassado, a mensagem é reentregue após um intervalo de nova tentativa. Defina esse valor com base na sua lógica de negócios.

15

--

Minutos

suspendTimeMillis

Intervalo de nova tentativa para mensagens ordenadas cujo consumo falhou.

--

--

--

maxCachedMessageAmount

Quantidade máxima de mensagens armazenadas em cache no cliente consumidor local. A cota é dividida igualmente entre todos os tópicos assinados. Por exemplo, se um consumidor assinar 2 tópicos e esse valor for 1.000, cada tópico poderá armazenar até 500 mensagens em cache. Caso o cliente consumidor busque várias mensagens por vez, o número real de mensagens em cache pode exceder esse valor. Recomenda-se definir esse parâmetro com aproximadamente o dobro da quantidade de mensagens que seu consumidor processa por segundo.

5.000

100--50.000

--

maxCachedMessageSizeInMiB

Tamanho total máximo das mensagens armazenadas em cache no cliente consumidor local.

512

16--2.048

MiB

Importante

Definir maxCachedMessageAmount ou maxCachedMessageSizeInMiB com valores muito altos pode causar erros de falta de memória (OOM) no cliente.

Parâmetros de push em lote

Estes parâmetros controlam como o broker agrupa mensagens no modo push antes de entregá-las.

Parâmetro

Descrição

Padrão

Valores válidos

Unidade

ConsumeMessageBatchMaxSize

Quantidade máxima de mensagens em cache entregues a um consumidor em um único lote. Quando o cache atinge esse limiar, todas as mensagens armazenadas são enviadas de uma só vez.

32

1--1.024

--

BatchConsumeMaxAwaitDurationInSeconds

Tempo máximo de espera antes que as mensagens em cache sejam enviadas aos consumidores de uma só vez.

0

0--450

Segundos

Parâmetros do modo pull

Os parâmetros a seguir aplicam-se exclusivamente a consumidores pull.

Parâmetro

Descrição

Padrão

Valores válidos

Unidade

maxCachedMessageSizeInMiB

Tamanho máximo das mensagens em cache por partição no cliente local.

100

16--2.048

MiB

autoCommit

Indica se os offsets do consumidor devem ser confirmados automaticamente.

true

true, false

--

autoCommitIntervalMillis

Intervalo entre confirmações automáticas de offset.

5

--

Segundos

pollTimeoutMillis

Tempo limite para cada chamada de poll().

5

--

Segundos

Importante

Configurar maxCachedMessageSizeInMiB com um valor excessivamente alto pode provocar erros de OOM no cliente.

Para mais informações sobre partições e offsets, consulte Termos.

Referências

Códigos de exemplo para envio e consumo de mensagens: