Documentation
¶
Index ¶
- Constants
- func IsTagType(exp string) bool
- type AllocateStrategy
- type ClientOption
- type ConsumeFromWhere
- type ConsumeResult
- type ConsumerOption
- type ExpressionType
- type Message
- type MessageExt
- type MessageModel
- type MessageQueue
- func AllocateByAveragely(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- func AllocateByAveragelyCircle(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- func AllocateByConfig(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- func AllocateByConsistentHash(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- func AllocateByMachineNearby(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- func AllocateByMachineRoom(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
- type MessageSelector
- type ProducerOptions
- type PullResult
- type PullStatus
- type SendResult
- type SendStatus
Constants ¶
const ( /** * <ul> * Keywords: * <li>{@code AND, OR, NOT, BETWEEN, IN, TRUE, FALSE, IS, NULL}</li> * </ul> * <p/> * <ul> * Data type: * <li>Boolean, like: TRUE, FALSE</li> * <li>String, like: 'abc'</li> * <li>Decimal, like: 123</li> * <li>Float number, like: 3.1415</li> * </ul> * <p/> * <ul> * Grammar: * <li>{@code AND, OR}</li> * <li>{@code >, >=, <, <=, =}</li> * <li>{@code BETWEEN A AND B}, equals to {@code >=A AND <=B}</li> * <li>{@code NOT BETWEEN A AND B}, equals to {@code >B OR <A}</li> * <li>{@code IN ('a', 'b')}, equals to {@code ='a' OR ='b'}, this operation only support String type.</li> * <li>{@code IS NULL}, {@code IS NOT NULL}, check parameter whether is null, or not.</li> * <li>{@code =TRUE}, {@code =FALSE}, check parameter whether is true, or false.</li> * </ul> * <p/> * <p> * Example: * (a > 10 AND a < 100) OR (b IS NOT NULL AND b=TRUE) * </p> */ SQL92 = ExpressionType("SQL92") /** * Only support or operation such as * "tag1 || tag2 || tag3", <br> * If null or * expression, meaning subscribe all. */ TAG = ExpressionType("TAG") )
const ( PropertyKeySeparator = " " PropertyKeys = "KEYS" PropertyTags = "TAGS" PropertyWaitStoreMsgOk = "WAIT" PropertyDelayTimeLevel = "DELAY" PropertyRetryTopic = "RETRY_TOPIC" PropertyRealTopic = "REAL_TOPIC" PropertyRealQueueId = "REAL_QID" PropertyTransactionPrepared = "TRAN_MSG" PropertyProducerGroup = "PGROUP" PropertyMinOffset = "MIN_OFFSET" PropertyMaxOffset = "MAX_OFFSET" PropertyBuyerId = "BUYER_ID" PropertyOriginMessageId = "ORIGIN_MESSAGE_ID" PropertyTransferFlag = "TRANSFER_FLAG" PropertyCorrectionFlag = "CORRECTION_FLAG" PropertyMQ2Flag = "MQ2_FLAG" PropertyReconsumeTime = "RECONSUME_TIME" PropertyMsgRegion = "MSG_REGION" PropertyTraceSwitch = "TRACE_ON" PropertyUniqueClientMessageIdKeyIndex = "UNIQ_KEY" PropertyMaxReconsumeTimes = "MAX_RECONSUME_TIMES" PropertyConsumeStartTime = "CONSUME_START_TIME" PropertyTranscationPreparedQueueOffset = "TRAN_PREPARED_QUEUE_OFFSET" PropertyTranscationCheckTimes = "TRANSACTION_CHECK_TIMES" PropertyCheckImmunityTimeInSeconds = "CHECK_IMMUNITY_TIME_IN_SECONDS" )
Variables ¶
This section is empty.
Functions ¶
Types ¶
type AllocateStrategy ¶
type AllocateStrategy func(string, string, []*MessageQueue, []string) []*MessageQueue
type ClientOption ¶
type ClientOption struct {
NameServerAddr string
ClientIP string
InstanceName string
UnitMode bool
UnitName string
VIPChannelEnabled bool
UseTLS bool
}
func (*ClientOption) ChangeInstanceNameToPID ¶
func (opt *ClientOption) ChangeInstanceNameToPID()
func (*ClientOption) String ¶
func (opt *ClientOption) String() string
type ConsumeFromWhere ¶
type ConsumeFromWhere int
Consuming point on consumer booting. </p>
There are three consuming points: <ul> <li> <code>CONSUME_FROM_LAST_OFFSET</code>: consumer clients pick up where it stopped previously. If it were a newly booting up consumer client, according aging of the consumer group, there are two cases: <ol> <li> if the consumer group is created so recently that the earliest message being subscribed has yet expired, which means the consumer group represents a lately launched business, consuming will start from the very beginning; </li> <li> if the earliest message being subscribed has expired, consuming will start from the latest messages, meaning messages born prior to the booting timestamp would be ignored. </li> </ol> </li> <li> <code>CONSUME_FROM_FIRST_OFFSET</code>: Consumer client will start from earliest messages available. </li> <li> <code>CONSUME_FROM_TIMESTAMP</code>: Consumer client will start from specified timestamp, which means messages born prior to {@link #consumeTimestamp} will be ignored </li> </ul>
const ( ConsumeFromLastOffset ConsumeFromWhere = iota ConsumeFromFirstOffset ConsumeFromTimestamp )
type ConsumeResult ¶
type ConsumeResult int
const ( ConsumeSuccess ConsumeResult = iota ConsumeRetryLater )
type ConsumerOption ¶
type ConsumerOption struct {
ClientOption
NameServerAddr string
/**
* Backtracking consumption time with second precision. Time format is
* 20131223171201<br>
* Implying Seventeen twelve and 01 seconds on December 23, 2013 year<br>
* Default backtracking consumption time Half an hour ago.
*/
ConsumeTimestamp string
// The socket timeout in milliseconds
ConsumerPullTimeout time.Duration
// Concurrently max span offset.it has no effect on sequential consumption
ConsumeConcurrentlyMaxSpan int
// Flow control threshold on queue level, each message queue will cache at most 1000 messages by default,
// Consider the {PullBatchSize}, the instantaneous value may exceed the limit
PullThresholdForQueue int64
// Limit the cached message size on queue level, each message queue will cache at most 100 MiB messages by default,
// Consider the {@code pullBatchSize}, the instantaneous value may exceed the limit
//
// The size of a message only measured by message body, so it's not accurate
PullThresholdSizeForQueue int
// Flow control threshold on topic level, default value is -1(Unlimited)
//
// The value of {@code pullThresholdForQueue} will be overwrote and calculated based on
// {@code pullThresholdForTopic} if it is't unlimited
//
// For example, if the value of pullThresholdForTopic is 1000 and 10 message queues are assigned to this consumer,
// then pullThresholdForQueue will be set to 100
PullThresholdForTopic int
// Limit the cached message size on topic level, default value is -1 MiB(Unlimited)
//
// The value of {@code pullThresholdSizeForQueue} will be overwrote and calculated based on
// {@code pullThresholdSizeForTopic} if it is't unlimited
//
// For example, if the value of pullThresholdSizeForTopic is 1000 MiB and 10 message queues are
// assigned to this consumer, then pullThresholdSizeForQueue will be set to 100 MiB
PullThresholdSizeForTopic int
// Message pull Interval
PullInterval time.Duration
// Batch consumption size
ConsumeMessageBatchMaxSize int
// Batch pull size
PullBatchSize int32
// Whether update subscription relationship when every pull
PostSubscriptionWhenPull bool
// Max re-consume times. -1 means 16 times.
//
// If messages are re-consumed more than {@link #maxReconsumeTimes} before success, it's be directed to a deletion
// queue waiting.
MaxReconsumeTimes int
// Suspending pulling time for cases requiring slow pulling like flow-control scenario.
SuspendCurrentQueueTimeMillis time.Duration
// Maximum amount of time a message may block the consuming thread.
ConsumeTimeout time.Duration
ConsumerModel MessageModel
Strategy AllocateStrategy
ConsumeOrderly bool
FromWhere ConsumeFromWhere
}
type ExpressionType ¶
type ExpressionType string
type Message ¶
type Message struct {
Topic string
Body []byte
Flag int32
Properties map[string]string
TransactionId string
Batch bool
}
func NewMessage ¶
func (*Message) PutProperty ¶
func (*Message) RemoveProperty ¶
type MessageExt ¶
type MessageExt struct {
Message
MsgId string
QueueId int32
StoreSize int32
QueueOffset int64
SysFlag int32
BornTimestamp int64
BornHost string
StoreTimestamp int64
StoreHost string
CommitLogOffset int64
BodyCRC int32
ReconsumeTimes int32
PreparedTransactionOffset int64
}
func (*MessageExt) GetTags ¶
func (msgExt *MessageExt) GetTags() string
func (*MessageExt) String ¶
func (msgExt *MessageExt) String() string
type MessageModel ¶
type MessageModel int
Message model defines the way how messages are delivered to each consumer clients. </p>
RocketMQ supports two message models: clustering and broadcasting. If clustering is set, consumer clients with the same {@link #consumerGroup} would only consume shards of the messages subscribed, which achieves load balances; Conversely, if the broadcasting is set, each consumer client will consume all subscribed messages separately. </p>
This field defaults to clustering.
const ( BroadCasting MessageModel = iota Clustering )
func (MessageModel) String ¶
func (mode MessageModel) String() string
type MessageQueue ¶
type MessageQueue struct {
Topic string `json:"topic"`
BrokerName string `json:"brokerName"`
QueueId int `json:"queueId"`
}
MessageQueue message queue
func AllocateByAveragely ¶
func AllocateByAveragely(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
func AllocateByAveragelyCircle ¶
func AllocateByAveragelyCircle(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
func AllocateByConfig ¶
func AllocateByConfig(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
func AllocateByConsistentHash ¶
func AllocateByConsistentHash(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
func AllocateByMachineNearby ¶
func AllocateByMachineNearby(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
TODO
func AllocateByMachineRoom ¶
func AllocateByMachineRoom(consumerGroup, currentCID string, mqAll []*MessageQueue, cidAll []string) []*MessageQueue
func (*MessageQueue) Equals ¶
func (mq *MessageQueue) Equals(queue *MessageQueue) bool
func (*MessageQueue) HashCode ¶
func (mq *MessageQueue) HashCode() int
func (*MessageQueue) String ¶
func (mq *MessageQueue) String() string
type MessageSelector ¶
type MessageSelector struct {
Type ExpressionType
Expression string
}
type ProducerOptions ¶
type ProducerOptions struct {
ClientOption
NameServerAddr string
GroupName string
RetryTimesWhenSendFailed int
UnitMode bool
}
type PullResult ¶
type PullResult struct {
NextBeginOffset int64
MinOffset int64
MaxOffset int64
Status PullStatus
SuggestWhichBrokerId int64
// contains filtered or unexported fields
}
PullResult the pull result
func (*PullResult) GetMessageExts ¶
func (result *PullResult) GetMessageExts() []*MessageExt
func (*PullResult) GetMessages ¶
func (result *PullResult) GetMessages() []*Message
func (*PullResult) SetMessageExts ¶
func (result *PullResult) SetMessageExts(msgExts []*MessageExt)
func (*PullResult) String ¶
func (result *PullResult) String() string
type PullStatus ¶
type PullStatus int
PullStatus pull status
const ( PullFound PullStatus = iota PullNoNewMsg PullNoMsgMatched PullOffsetIllegal PullBrokerTimeout )
predefined pull status
type SendResult ¶
type SendResult struct {
Status SendStatus
MsgID string
MessageQueue *MessageQueue
QueueOffset int64
TransactionID string
OffsetMsgID string
RegionID string
TraceOn bool
}
SendResult RocketMQ send result
func (*SendResult) String ¶
func (result *SendResult) String() string
SendResult send message result to string(detail result)
type SendStatus ¶
type SendStatus int
SendStatus of message
const ( SendOK SendStatus = iota SendFlushDiskTimeout SendFlushSlaveTimeout SendSlaveNotAvailable FlagCompressed = 0x1 MsgIdLength = 8 + 8 )