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)
}
必要に応じて concurrency や connections を定義できる。下にある 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]}
すると、stdout に worker と [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 は主に以下の処理を行っている:
poller.getJob()を介してLPOPを呼び出し、タスクを取り出す
reply, err := conn.Do("LPOP", fmt.Sprintf("%squeue:%s", workerSettings.Namespace, queue))
poller.poll(duration, quit)を介して定期的に jobs チャンネルから実行対象のjobsを取り出し、quitを終了シグナルとして利用する。
jobs := make(chan *Job)
//......
return jobs
- 最後に
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 の使用率が目に見えて跳ね上がっているわけでもないのだが……。
関連記事
- 測定が目標になるとき:窓税からPull Request数まで かつて僕は小さなツールを自作し、四半期で自分がどれだけPRに貢献したか、レビューコメントをどれだけ残したか、チケットをどれだけ消化したかを集計して、上司にアウトプットを証明しようとしたことがある。上司は淡々と、評価はアウトプットだけで見るものではないと言った。数年後、僕はようやく理解した――測定が目標になるとき、それはもはや良い測定ではなくなるのだ。英国の窓税、ハノイのネズミ駆除の報奨金から、現代のPR数による開発者評価に至るまで、そのメカニズムはまったく同じだ。
- Cloudflare Images を画像ストレージ・変換ソリューションとして使う ウェブページに画像を1枚置くのはフロントエンドにとって最も簡単なことだが、リサイズや各種フォーマットの生成、さらにはトラフィックの負荷に耐えることまで完璧にやろうとすると、実際には一つの包括的なソリューションが必要になる。僕はその後、すべて Cloudflare Images に任せるようになり、オリジナル画像1枚だけを渡すようにしている。
- もう AWS Access Key を使うのはやめよう Access Key は AWS において見落とされがちなセキュリティリスクだ。OIDC と IAM Role を組み合わせることで、GitHub Actions にシークレットを一切保持させることなく、安全に AWS リソースを操作できるようにする。
- データベース主キー:AUTO_INCREMENT、UUID、そしてUUIDv7 バックエンド開発で度々直面する主キーの決定。auto incrementを使うべきか、それともUUIDか?衝突への懸念は?UUIDv7とcreated_at + インデックスの性能差はどれほどか?実際に2,000万件のデータで検証したベンチマークと設計上の意思決定を解説する。