【问题标题】:process one message per user每个用户处理一条消息
【发布时间】:2020-08-31 03:17:06
【问题描述】:

我在 redis 中有一个列表,用作队列。我在左侧推送元素并从右侧弹出。来自不同用户的请求被推送到队列中。我有一个 goroutines 池,它们从队列(POP)中读取请求并处理它们。我希望一次只能处理每个 userId 的一个请求。我有一个永远运行的 ReadRequest() 函数,它会弹出一个具有 userId 的请求。我需要按用户进入的顺序处理每个用户的请求。我不确定如何实现这一点。我需要每个 userId 的 redis 列表吗?如果是这样,我将如何遍历处理其中请求的所有列表?

for i:=0; i< 5; i++{
  wg.Add(1)
  go ReadRequest(&wg)

}


func ReadRequest(){

   for{

      //redis pop request off list
       request:=MyRedisPop()
       fmt.Println(request.UserId)

      // only call Process if no other goroutine is processing a request for this user
      Process(request)



 time.sleep(100000)
     }

wg.Done()

}

【问题讨论】:

  • 是的,您希望每个用户有一个列表。您希望每个列表有一个 goroutine,每个列表轮询一个列表。
  • @Adrian 如何为每个用户队列启动一个 goroutine
  • 我不确定我是否理解你在问什么 - 你在引用的代码中表明你知道如何启动一个 goroutine。
  • 如果一个新用户连接并在redis中创建一个新列表,go代码如何知道开始一个新的goroutine?
  • 如果列表不是由同一个应用程序创建的,您可以定期轮询 redis 以检查新列表。

标签: go redis queue message-queue


【解决方案1】:

以下是无需创建多个 Redis 列表即可使用的伪代码:

// maintain a global map for all users
// if you see a new user, call NewPerUser() and add it to the list
// Then, send the request to the corresponding channel for processing
var userMap map[string]PerUser 

type PerUser struct {
    chan<- redis.Request // Whatever is the request type
    semaphore *semaphore.Weighted // Semaphore to limit concurrent processing
}

func NewPerUser() *PerUser {
    ch := make(chan redis.Request)
    s := semaphore.NewWeighted(1) // One 1 concurrent request is allowed
    go func(){
        for req := range ch {
            s.Acquire(context.Background(), 1)
            defer s.Release(1)
            // Process the request here
        }
    }()
}

请注意,这只是一个伪代码,我还没有测试它是否有效。

【讨论】:

  • 如何在用户不再连接后停止 goroutine
  • 你可以关闭相应的通道——它将退出for循环并结束go-routine。
猜你喜欢
  • 1970-01-01
  • 2018-12-28
  • 1970-01-01
  • 2020-08-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-26
  • 1970-01-01
相关资源
最近更新 更多