2
0
mirror of https://github.com/hibiken/asynq.git synced 2025-10-25 10:56:12 +08:00

Add support for multiple queues

This commit is contained in:
Ken Hibino
2019-11-16 08:20:23 -08:00
parent 95023bd3b5
commit f4d59bece7

View File

@@ -100,11 +100,11 @@ func (w *Workers) Run() {
for { for {
// pull message out of the queue and process it // pull message out of the queue and process it
// TODO(hibiken): get a list of queues in order of priority // TODO(hibiken): sort the list of queues in order of priority
res, err := w.rdb.BLPop(0, "asynq:queues:test").Result() // A timeout of zero means block indefinitely. res, err := w.rdb.BLPop(0, listQueues(w.rdb)...).Result() // A timeout of zero means block indefinitely.
if err != nil { if err != nil {
if err != redis.Nil { if err != redis.Nil {
log.Printf("error when BLPOP from %s: %v\n", "aysnq:queues:test", err) log.Printf("BLPOP command failed: %v\n", err)
} }
continue continue
} }
@@ -163,15 +163,22 @@ func (w *Workers) pollScheduledTasks() {
} }
} }
// push pushes the task to the specified queue to get picked up by a worker.
func push(rdb *redis.Client, queue string, t *Task) error { func push(rdb *redis.Client, queue string, t *Task) error {
bytes, err := json.Marshal(t) bytes, err := json.Marshal(t)
if err != nil { if err != nil {
return fmt.Errorf("could not encode task into JSON: %v", err) return fmt.Errorf("could not encode task into JSON: %v", err)
} }
err = rdb.SAdd(allQueues, queue).Err() qname := queuePrefix + queue
err = rdb.SAdd(allQueues, qname).Err()
if err != nil { if err != nil {
return fmt.Errorf("could not execute command SADD %q %q: %v", return fmt.Errorf("could not execute command SADD %q %q: %v",
allQueues, queue, err) allQueues, qname, err)
} }
return rdb.RPush(queuePrefix+queue, string(bytes)).Err() return rdb.RPush(qname, string(bytes)).Err()
}
// listQueues returns the list of all queues.
func listQueues(rdb *redis.Client) []string {
return rdb.SMembers(allQueues).Val()
} }