· 4 min read

Introduction to Goworker — Implementing Workers with Redis

This article was auto-translated from Chinese. Some nuances may be lost in translation.

Recently, I’ve been gradually migrating a project from Ruby on Rails to Golang. Part of the reason is wanting to practice writing its rather quirky syntax to see how it feels, and another part is experiencing Golang’s power firsthand while learning how to organize code without an established framework.

Today, I’m going to introduce goworker. The reason it caught my eye is that it was the first result when searching for “go worker,” and its implementation is fully compatible with the format used by resque (a Ruby worker gem). This means I can push tasks into a queue from a Ruby on Rails application and have them processed via Golang.

How to Use It

During initialization, you need to register all workers that will be executed with goworker so that it knows which corresponding task to run when pulling from a given queue:

func init() {
  settings := goworker.WorkerSettings{
    URI:            os.Getenv("REDIS_URL"),
    Queues:         []string{"worker", "queues"},
    UseNumber:      true,
    ExitOnComplete: false, // don't complete even though queue is empty.
    Concurrency:    concurrency,
    Connections:    connections,
    Namespace:      "myapp" + ":",
    Interval:       10.0,
  }
  goworker.SetSettings(settings)  
}

You can define concurrency and connections according to your needs. The Namespace below provides a scope for the queue names. Interval specifies how long to wait before checking for new tasks again when the queue is empty.

How do you implement a worker? You need to implement a function and register it with goworker.

func workerFn(queue string, args ...interface{}) error {
  // your job.
  fmt.Println(queue, args)
}

goworker.Register("MyWorker", workerFn)

Once configured, you can use goworker.Work() to listen to Redis. This is a blocking operation, so it’s typically run on a separate server or process (you can also spawn a goroutine for it when starting the server).

func main() {
	err := goworker.Work()
  if err != nil {
    // log your error
  }  
}

Next, let’s push a task into Redis to see it in action:

RPUSH myapp:queue:worker '{"class": "MyWorker", "args": [1,2,3]}

You’ll see worker and [1 2 3] printed to stdout.

You can also push tasks directly into Redis using goworker.Enqueue(&goworker.Job{}).

How It Works

When a new task is added, goworker pushes it to namespace:queue:job.Queue using RPUSH. To serialize the parameters, a JSON string is used as the payload and pushed into the queue.

// workers.go L43~47
buffer, err := json.Marshal(job.Payload)
	if err != nil {
		logger.Criticalf("Cant marshal payload on enqueue")
		return err
	}

	err = conn.Send("RPUSH", fmt.Sprintf("%squeue:%s", workerSettings.Namespace, job.Queue), buffer)
	if err != nil {
		logger.Criticalf("Cant push to queue")
		return err
	}

In addition, a set is used to keep track of which queues are currently in use.

// workers.go L49~53
err = conn.Send("SADD", fmt.Sprintf("%squeues", workerSettings.Namespace), job.Queue)
if err != nil {
  logger.Criticalf("Cant register queue to list of use queues")
  return err
}

When calling goworker.Work(), it internally invokes a poller. This poller mainly performs a few tasks:

  1. Retrieves tasks by calling LPOP via poller.getJob()
reply, err := conn.Do("LPOP", fmt.Sprintf("%squeue:%s", workerSettings.Namespace, queue))
  1. Periodically retrieves jobs to be executed from a jobs channel via poller.poll(duration, quit), using quit as an exit signal.
jobs := make(chan *Job)
//......
return jobs
  1. Finally, executes worker.work(jobs, &monitor). When a job is received, it runs w.run(job, workerFunc), which is the function we initially defined. If a matching class cannot be found, it prints the error No Worker for ... queue with args ....

Caveats

  • Using goworker does not guarantee that pending jobs will run successfully if an unexpected shutdown occurs.
  • As far as I understand, Redis doesn’t have a built-in mechanism to guarantee that the receiving end will definitely receive the message. If you need stronger reliability guarantees, protocols like AMQP might be necessary. However, for a small project, Redis is more than enough.

Summary

For general applications, Redis is quite sufficient. Unless you are dealing with genuinely massive traffic and high-concurrency scenarios, Redis-based queues are far from inadequate. In this post, we looked at developing with the goworker library paired with Redis, which should be well-suited for handling time-consuming (or delayed) background tasks. For larger-scale applications, more reliable protocols like AMQP can be adopted.

Also, my computer fans kept spinning constantly whenever I ran goworker. I’m not sure if that’s normal, though the CPU usage didn’t show any obvious spikes…

Related Posts

Explore Other Topics