· 4分で読了

Goworker 入門 — Redis と組み合わせた worker の実装

この記事は中国語から自動翻訳されたものです。翻訳によりニュアンスが失われている場合があります。

最近、プロジェクトを Ruby on Rails から Go 言語(golang)へと徐々に移行している。少し不恰好な構文を実際に書いてみてどんな感触かを試したいというのもあるし、Go 言語の威力や、決まった既存フレームワークがない中でどのようにコードを構成していくかを体感してみたいという理由もある。

今回紹介するのは goworker だ。これに目をつけた理由は、「go worker」と検索して一番上に出てきたからであり、その実装方法が resque(Ruby のワーカーライブラリ)の形式を完全にサポートしているためだ。これにより、Ruby on Rails アプリケーション側でキューにタスクを投入し、Go 言語側でタスクを実行することができる。

使い方

初期化時には、実行されるすべてのワーカーをあらかじめ goworker に登録しておく必要がある。これにより、ワーカーがキューからジョブを取得した際に、どの対応するタスクを実行すべきかを認識できるようになる。

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)  
}

必要に応じて concurrencyconnections を定義できる。下にある Namespace はキューの名前にスコープ(プレフィックス)を持たせるためのものだ。interval は、キューにタスクがない場合に、どれくらいの間隔を空けて新しいタスクの有無を再確認するかを指定する。

ワーカーを実装するにはどうすればいいのか?関数を一つ実装し、その関数を goworker に登録すればよい。

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

goworker.Register("MyWorker", workerFn)

設定完了後は goworker.Start() を使って Redis をリッスンできる。この処理はブロッキングするため、通常は別サーバーまたは別プロセスとして立ち上げる(サーバー起動時に goroutine を立ち上げて実行しても構わない)。

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

次に、Redis にタスクを投入してみよう。

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

すると、stdoutworker[1 2 3] が出力されるのが確認できる。

また、直接 goworker.Enqueue(&goworker.Job{}) を使って Redis にタスクをプッシュすることも可能だ。

仕組み

新しいタスクが追加されると、goworker は RPUSH を使ってタスクを namespace:queue:job.Queue に投入する。パラメータをシリアライズできるようにするため、JSON 文字列をペイロードとしてキューに追加している。

// 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
	}

さらに、現在使用中のキューの一覧をセット(集合)を使って記録している。

// 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
}

goworker.Work() を呼び出すと、内部でまず poller が呼び出される。この poller は主に以下の処理を行っている:

  1. poller.getJob() を介して LPOP を呼び出し、タスクを取り出す
reply, err := conn.Do("LPOP", fmt.Sprintf("%squeue:%s", workerSettings.Namespace, queue))
  1. poller.poll(duration, quit) を介して定期的に jobs チャンネルから実行対象の jobs を取り出し、quit を終了シグナルとして利用する。
jobs := make(chan *Job)
//......
return jobs
  1. 最後に worker.work(jobs, &monitor) を実行し、ジョブを受け取ると w.run(job, workerFunc)(最初に定義した関数)を実行する。対応するクラスが見つからない場合は、No Worker for ... queue with args ... というエラーが出力される。

課題・懸念点

  • goworker を利用する場合、予期せずシャットダウンしたときに、保留中(pending)のジョブが正常に実行される保証はない。
  • 僕の理解では、Redis には受信側が確実にメッセージを受け取ることを保証する仕組みがないため、より高信頼な仕組みが必要であれば AMQP などのプロトコルを使う必要があるだろう。ただ、小規模なプロジェクトであれば Redis で十分事足りる。

まとめ

一般的なアプリケーションであれば Redis で十分に対応できる。よほどの大規模トラフィックや超高並行処理の現場でない限り、Redis のキューは決して使い物にならないわけではない。今回は goworker というライブラリと Redis を組み合わせた開発を紹介したが、時間のかかるタスク(または遅延タスク)の処理には十分対応できるはずだ。より大規模なアプリケーションの場合は、AMQP などのより信頼性の高いプロトコルを採用すると良いだろう。 また、goworker を実行していると僕の PC のファンがずっと回りっぱなしになるのだが、これが正常なのかどうかは分からない。ただ、CPU の使用率が目に見えて跳ね上がっているわけでもないのだが……。

関連記事

他のトピックを探索