@ -77,27 +77,13 @@ func (mc *OnlineHistoryMongoConsumerHandler) handleChatWs2Mongo(ctx context.Cont
for _ , msg := range msgFromMQ . MsgData {
seqs = append ( seqs , msg . Seq )
}
err = mc . msgTransferDatabase . DeleteMessagesFromCache ( ctx , msgFromMQ . ConversationID , seqs )
if err != nil {
log . ZError (
ctx ,
"remove cache msg from redis err" ,
err ,
"msg" ,
msgFromMQ . MsgData ,
"conversationID" ,
msgFromMQ . ConversationID ,
)
}
}
func ( * OnlineHistoryMongoConsumerHandler ) Setup ( _ sarama . ConsumerGroupSession ) error { return nil }
func ( * OnlineHistoryMongoConsumerHandler ) Cleanup ( _ sarama . ConsumerGroupSession ) error { return nil }
func ( mc * OnlineHistoryMongoConsumerHandler ) ConsumeClaim (
sess sarama . ConsumerGroupSession ,
claim sarama . ConsumerGroupClaim ,
) error { // an instance in the consumer group
func ( mc * OnlineHistoryMongoConsumerHandler ) ConsumeClaim ( sess sarama . ConsumerGroupSession , claim sarama . ConsumerGroupClaim ) error { // an instance in the consumer group
log . ZDebug ( context . Background ( ) , "online new session msg come" , "highWaterMarkOffset" ,
claim . HighWaterMarkOffset ( ) , "topic" , claim . Topic ( ) , "partition" , claim . Partition ( ) )
for msg := range claim . Messages ( ) {