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.
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 |
|
|
Endpoint TCP da instância do ApsaraMQ for RocketMQ. Encontre essa informação na página Instance Details no console. |
-- |
|
|
AccessKey ID usado como identificador exclusivo para autenticação. Consulte Crie um par de AccessKey. |
-- |
|
|
Segredo do AccessKey usado como senha para autenticação. Consulte Crie um par de AccessKey. |
-- |
|
|
Canal do usuário. Defina como |
|
Métodos do produtor
A interface do produtor fornece os seguintes métodos:

Parâmetros do produtor
|
Parâmetro |
Descrição |
Padrão |
Unidade |
|
|
Tempo limite para envio de uma mensagem. |
-- |
Milissegundos |
|
|
Tempo mínimo de espera antes que o broker verifique pela primeira vez o status de uma mensagem transacional. |
-- |
Segundos |
|
|
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:

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

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 |
|
|
ID do grupo de consumidores criado no console do ApsaraMQ for RocketMQ. Consulte Termos. |
-- |
-- |
-- |
|
|
Modelo de consumo. Valores válidos: |
|
|
-- |
|
|
Quantidade de threads que o consumidor utiliza para processar mensagens. |
20 |
-- |
-- |
|
|
Número máximo de tentativas de reconsumo quando falha o consumo de uma mensagem. |
16 |
-- |
-- |
|
|
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 |
|
|
Intervalo de nova tentativa para mensagens ordenadas cujo consumo falhou. |
-- |
-- |
-- |
|
|
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 |
-- |
|
|
Tamanho total máximo das mensagens armazenadas em cache no cliente consumidor local. |
512 |
16--2.048 |
MiB |
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 |
|
|
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 |
-- |
|
|
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 |
|
|
Tamanho máximo das mensagens em cache por partição no cliente local. |
100 |
16--2.048 |
MiB |
|
|
Indica se os offsets do consumidor devem ser confirmados automaticamente. |
|
|
-- |
|
|
Intervalo entre confirmações automáticas de offset. |
5 |
-- |
Segundos |
|
|
Tempo limite para cada chamada de |
5 |
-- |
Segundos |
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: