mirror of
https://github.com/openimsdk/open-im-server.git
synced 2025-11-05 03:42:08 +08:00
feat:optimize GetConversationsHasReadAndMaxSeq
This commit is contained in:
parent
fe7c029c2a
commit
04a5c97e31
19
pkg/common/storage/cache/redis/seq.go
vendored
19
pkg/common/storage/cache/redis/seq.go
vendored
@ -16,6 +16,7 @@ package redis
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
||||||
"github.com/openimsdk/tools/errs"
|
"github.com/openimsdk/tools/errs"
|
||||||
@ -61,17 +62,27 @@ func (c *seqCache) getSeq(ctx context.Context, conversationID string, getkey fun
|
|||||||
|
|
||||||
func (c *seqCache) getSeqs(ctx context.Context, items []string, getkey func(s string) string) (m map[string]int64, err error) {
|
func (c *seqCache) getSeqs(ctx context.Context, items []string, getkey func(s string) string) (m map[string]int64, err error) {
|
||||||
m = make(map[string]int64, len(items))
|
m = make(map[string]int64, len(items))
|
||||||
|
keys := make([]string, len(items))
|
||||||
for i, v := range items {
|
for i, v := range items {
|
||||||
res, err := c.rdb.Get(ctx, getkey(v)).Result()
|
keys[i] = getkey(v)
|
||||||
if err != nil && err != redis.Nil {
|
}
|
||||||
|
|
||||||
|
res, err := c.rdb.MGet(ctx, keys...).Result()
|
||||||
|
if err != nil && !errors.Is(err, redis.Nil) {
|
||||||
return nil, errs.Wrap(err)
|
return nil, errs.Wrap(err)
|
||||||
}
|
}
|
||||||
val := stringutil.StringToInt64(res)
|
|
||||||
|
// len(res) == len(items)
|
||||||
|
for i := range res {
|
||||||
|
strRes, ok := res[i].(string)
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
val := stringutil.StringToInt64(strRes)
|
||||||
if val != 0 {
|
if val != 0 {
|
||||||
m[items[i]] = val
|
m[items[i]] = val
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return m, nil
|
return m, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -16,15 +16,21 @@ package rpccache
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/localcache"
|
"github.com/openimsdk/open-im-server/v3/pkg/localcache"
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
||||||
pbconversation "github.com/openimsdk/protocol/conversation"
|
pbconversation "github.com/openimsdk/protocol/conversation"
|
||||||
"github.com/openimsdk/tools/errs"
|
"github.com/openimsdk/tools/errs"
|
||||||
"github.com/openimsdk/tools/log"
|
"github.com/openimsdk/tools/log"
|
||||||
|
"github.com/openimsdk/tools/mq/memamq"
|
||||||
"github.com/redis/go-redis/v9"
|
"github.com/redis/go-redis/v9"
|
||||||
|
"sync"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
notificationWorkerCount = 2
|
||||||
|
notificationBufferSize = 200
|
||||||
)
|
)
|
||||||
|
|
||||||
func NewConversationLocalCache(client rpcclient.ConversationRpcClient, localCache *config.LocalCache, cli redis.UniversalClient) *ConversationLocalCache {
|
func NewConversationLocalCache(client rpcclient.ConversationRpcClient, localCache *config.LocalCache, cli redis.UniversalClient) *ConversationLocalCache {
|
||||||
@ -39,6 +45,7 @@ func NewConversationLocalCache(client rpcclient.ConversationRpcClient, localCach
|
|||||||
localcache.WithLocalSuccessTTL(lc.Success()),
|
localcache.WithLocalSuccessTTL(lc.Success()),
|
||||||
localcache.WithLocalFailedTTL(lc.Failed()),
|
localcache.WithLocalFailedTTL(lc.Failed()),
|
||||||
),
|
),
|
||||||
|
queue: memamq.NewMemoryQueue(notificationWorkerCount, notificationBufferSize),
|
||||||
}
|
}
|
||||||
if lc.Enable() {
|
if lc.Enable() {
|
||||||
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
||||||
@ -49,6 +56,8 @@ func NewConversationLocalCache(client rpcclient.ConversationRpcClient, localCach
|
|||||||
type ConversationLocalCache struct {
|
type ConversationLocalCache struct {
|
||||||
client rpcclient.ConversationRpcClient
|
client rpcclient.ConversationRpcClient
|
||||||
local localcache.Cache[any]
|
local localcache.Cache[any]
|
||||||
|
|
||||||
|
queue *memamq.MemoryQueue
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *ConversationLocalCache) GetConversationIDs(ctx context.Context, ownerUserID string) (val []string, err error) {
|
func (c *ConversationLocalCache) GetConversationIDs(ctx context.Context, ownerUserID string) (val []string, err error) {
|
||||||
@ -91,16 +100,43 @@ func (c *ConversationLocalCache) GetSingleConversationRecvMsgOpt(ctx context.Con
|
|||||||
|
|
||||||
func (c *ConversationLocalCache) GetConversations(ctx context.Context, ownerUserID string, conversationIDs []string) ([]*pbconversation.Conversation, error) {
|
func (c *ConversationLocalCache) GetConversations(ctx context.Context, ownerUserID string, conversationIDs []string) ([]*pbconversation.Conversation, error) {
|
||||||
conversations := make([]*pbconversation.Conversation, 0, len(conversationIDs))
|
conversations := make([]*pbconversation.Conversation, 0, len(conversationIDs))
|
||||||
|
|
||||||
|
errChan := make(chan error, len(conversationIDs))
|
||||||
|
conversationsChan := make(chan *pbconversation.Conversation, len(conversationIDs))
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
wg.Add(len(conversationIDs))
|
||||||
|
|
||||||
for _, conversationID := range conversationIDs {
|
for _, conversationID := range conversationIDs {
|
||||||
|
err := c.queue.Push(func() {
|
||||||
|
defer wg.Done()
|
||||||
conversation, err := c.GetConversation(ctx, ownerUserID, conversationID)
|
conversation, err := c.GetConversation(ctx, ownerUserID, conversationID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if errs.ErrRecordNotFound.Is(err) {
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
continue
|
return
|
||||||
}
|
}
|
||||||
|
errChan <- err
|
||||||
|
return
|
||||||
|
}
|
||||||
|
conversationsChan <- conversation
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
// push err
|
||||||
|
return nil, errs.Wrap(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
close(errChan)
|
||||||
|
close(conversationsChan)
|
||||||
|
|
||||||
|
err, ok := <-errChan
|
||||||
|
if ok {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
for conversation := range conversationsChan {
|
||||||
conversations = append(conversations, conversation)
|
conversations = append(conversations, conversation)
|
||||||
}
|
}
|
||||||
|
|
||||||
return conversations, nil
|
return conversations, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user