RocketMQ Configuration
Reference · Applies to Brighter V10 · Prerequisites: Basic Configuration
Apache RocketMQ is a distributed messaging platform, and Brighter configures it with a gateway connection wrapping the RocketMQ client, a publication per topic, and a subscription per consumer group.
RocketMQ General
Install the transport package:
dotnet add package Paramore.Brighter.MessagingGateway.RocketMQBrighter publishes to a RocketMQ topic and consumes with a SimpleConsumer in a consumer group, filtering by tag. Three things shape how you configure it:
Brighter cannot create a RocketMQ topic. The RocketMQ C# client has no administrative API, so a publication with
MakeChannelsset toCreatelogs a warning and carries on. Provision topics and consumer groups with the RocketMQ tooling, and setmakeChannelstoOnMissingChannel.Assume.The consumer group is a subscription option, and it has no usable default. It defaults to an empty string, which RocketMQ rejects, so set
consumerGroupon every subscription.The endpoint lives on the RocketMQ client's own configuration, not on a Brighter option.
RocketMessagingGatewayConnectiontakes anOrg.Apache.Rocketmq.ClientConfigas its constructor argument and exposes it as the get-onlyClientConfigproperty, so the endpoint, TLS and request timeout are set withClientConfig.Builder()and are not in the table below.
The consumer requires a RocketSubscription; a base Subscription raises a ConfigurationException when the channel is created.
RocketMQ Connection
RocketMessagingGatewayConnection takes its options as properties, alongside the ClientConfig it is constructed with.
TimerProvider
TimeProvider
TimeProvider.System
The time source used for delayed and time-dependent operations.
MaxAttempts
int
3
Attempts the producer makes to deliver a message.
Checker
ITransactionChecker?
null
The callback RocketMQ uses to resolve the state of a half-committed transactional message.
Instrumentation
InstrumentationOptions
All
The telemetry detail producers and consumers emit.
The property is spelled TimerProvider, with an r, where the rest of Brighter spells it TimeProvider.
RocketMQ Publication
RocketMqPublication takes its options as properties and adds these three to the base publication options, which it inherits.
Tag
string?
null
The tag messages are published with, which consumers filter on.
Instrumentation
InstrumentationOptions?
null
The telemetry detail this publication emits; null uses the connection's setting.
TopicType
TopicType
Normal
Selects a normal, delay or FIFO topic.
RocketMQ Subscription
The non-generic subscription type is RocketSubscription, and it takes its options as constructor arguments, so the option is the parameter you type. The seventeen it shares with Subscription behave the same way here; the other six are RocketMQ's own, plus Brighter's dead letter and invalid message routing keys.
subscriptionName
SubscriptionName
none
Names the subscription for diagnostics; read back as Name.
channelName
ChannelName
none
Names the channel this subscription reads.
routingKey
RoutingKey
none
The RocketMQ topic the consumer subscribes to.
requestType
Type?
none
The request type messages on this topic are translated into.
getRequestType
Func<Message, Type>?
derives the type from requestType
Determines the request type from the message rather than from the topic.
consumerGroup
string?
""
The RocketMQ consumer group this consumer joins.
bufferSize
int
1
Messages received at once and held in the channel.
noOfPerformers
int
1
Threads reading this topic, each with its own message pump.
timeOut
TimeSpan?
300 ms
How long a read waits before treating the channel as empty.
requeueCount
int
-1
Times a message is requeued before it is treated as a poison pill; -1 is unlimited.
requeueDelay
TimeSpan?
0 ms
How long delivery of a requeued message is delayed.
unacceptableMessageLimit
int
0
Unacceptable messages before the channel stops; 0 disables the limit.
unacceptableMessageLimitWindow
TimeSpan?
null
The window the unacceptable-message count resets at the end of.
messagePumpType
MessagePumpType
none
Selects the Reactor or Proactor concurrency model.
channelFactory
IAmAChannelFactory?
null
Creates the channel; supply a RocketMqChannelFactory over a RocketMessageConsumerFactory.
makeChannels
OnMissingChannel
Create
Whether Brighter creates missing infrastructure, validates it, or assumes it.
filter
FilterExpression?
every tag
The tag expression the consumer subscribes with.
emptyChannelDelay
TimeSpan?
500 ms
How long the pump pauses after a read that found no message.
channelFailureDelay
TimeSpan?
1000 ms
How long the pump pauses after a channel failure.
receiveMessageTimeout
TimeSpan?
60000 ms
How long a receive call waits at the broker before returning empty.
invisibilityTimeout
TimeSpan?
30000 ms
How long a received message is hidden from other consumers in the group.
deadLetterRoutingKey
RoutingKey?
null
The routing key messages are dead-lettered to.
invalidMessageRoutingKey
RoutingKey?
null
The routing key unacceptable messages are routed to.
The generic form is spelled differently from the non-generic one: RocketMqSubscription<T> derives from RocketSubscription. It supplies requestType from T and takes the same options otherwise. It still requires a subscription name, a channel name and a routing key, and it leaves messagePumpType required, so state Reactor or Proactor on every subscription.
RocketMQ Configuration Example
RocketMQ ships a message producer factory rather than a producer registry factory, so the registry is built from the factory's dictionary.
Further Reading
Basic Configuration — registering Brighter
Dispatcher Configuration Reference — the subscription options every transport shares
Command Processor Configuration Reference — the publication options every transport shares
Reactor and Proactor — choosing a message pump
Error Handling Options — what Brighter does with a message it cannot process
Last updated
Was this helpful?
