mirror of
https://github.com/openimsdk/open-im-server.git
synced 2025-04-25 11:06:43 +08:00
concurrent consumption of messages
This commit is contained in:
parent
930923e330
commit
15329a97fe
@ -96,8 +96,8 @@ func (och *OnlineHistoryConsumerHandler) Run(channelID int) {
|
|||||||
msgChannelValue := cmd.Value.(MsgChannelValue)
|
msgChannelValue := cmd.Value.(MsgChannelValue)
|
||||||
msgList := msgChannelValue.msgList
|
msgList := msgChannelValue.msgList
|
||||||
triggerID := msgChannelValue.triggerID
|
triggerID := msgChannelValue.triggerID
|
||||||
storageMsgList := make([]*pbMsg.MsgDataToMQ, 80)
|
storageMsgList := make([]*pbMsg.MsgDataToMQ, 0, 80)
|
||||||
pushMsgList := make([]*pbMsg.MsgDataToMQ, 80)
|
pushMsgList := make([]*pbMsg.MsgDataToMQ, 0, 80)
|
||||||
log.Debug(triggerID, "msg arrived channel", "channel id", channelID, msgList, msgChannelValue.userID, len(msgList))
|
log.Debug(triggerID, "msg arrived channel", "channel id", channelID, msgList, msgChannelValue.userID, len(msgList))
|
||||||
for _, v := range msgList {
|
for _, v := range msgList {
|
||||||
log.Debug(triggerID, "msg come to storage center", v.String())
|
log.Debug(triggerID, "msg come to storage center", v.String())
|
||||||
@ -119,7 +119,7 @@ func (och *OnlineHistoryConsumerHandler) Run(channelID int) {
|
|||||||
// log.NewError(msgFromMQ.OperationID, "SessionType error", msgFromMQ.String())
|
// log.NewError(msgFromMQ.OperationID, "SessionType error", msgFromMQ.String())
|
||||||
// return
|
// return
|
||||||
//}
|
//}
|
||||||
|
log.Debug(triggerID, "msg storage length", len(storageMsgList), "push length", len(pushMsgList))
|
||||||
err := saveUserChatList(msgChannelValue.userID, storageMsgList, triggerID)
|
err := saveUserChatList(msgChannelValue.userID, storageMsgList, triggerID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
singleMsgFailedCount += uint64(len(storageMsgList))
|
singleMsgFailedCount += uint64(len(storageMsgList))
|
||||||
@ -227,7 +227,7 @@ func (och *OnlineHistoryConsumerHandler) MessagesDistributionHandle() {
|
|||||||
log.Error(triggerID, "msg_transfer Unmarshal msg err", "msg", string(consumerMessages[i].Value), "err", err.Error())
|
log.Error(triggerID, "msg_transfer Unmarshal msg err", "msg", string(consumerMessages[i].Value), "err", err.Error())
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
log.Debug(triggerID, "single msg come to distribution center", msgFromMQ.String())
|
log.Debug(triggerID, "single msg come to distribution center", msgFromMQ.String(), string(consumerMessages[i].Key))
|
||||||
if oldM, ok := UserAggregationMsgs[string(consumerMessages[i].Key)]; ok {
|
if oldM, ok := UserAggregationMsgs[string(consumerMessages[i].Key)]; ok {
|
||||||
oldM = append(oldM, &msgFromMQ)
|
oldM = append(oldM, &msgFromMQ)
|
||||||
UserAggregationMsgs[string(consumerMessages[i].Key)] = oldM
|
UserAggregationMsgs[string(consumerMessages[i].Key)] = oldM
|
||||||
|
Loading…
x
Reference in New Issue
Block a user