ラベル celery の投稿を表示しています。 すべての投稿を表示
ラベル celery の投稿を表示しています。 すべての投稿を表示

2025年3月13日木曜日

OpenTelemetry で celery の非同期タスクでもリクエストと同じトレースをする方法

OpenTelemetry で celery の非同期タスクでもリクエストと同じトレースをする方法

概要

Flask アプリから celery の非同期タスクを呼び出した際にもトレースできるようにします
Flask、Celery で一貫して OLTP にトレース情報を送信するサンプルコードを紹介します

環境

  • macOS 15.3.1
  • Python 3.12.9
    • flask 3.1.0
    • open-telemetry 0.51b0
    • celery 5.4.0
  • OpenTelemetry Collector 0.121.0
  • Jaeger (all in one) 1.67.0

追加インストール

前回のに追加で以下をインストールします

  • pipenv install opentelemetry-instrumentation-celery celery

app.py

from celery import Celery
from flask import Flask
from opentelemetry import trace
from opentelemetry.propagate import inject

from lib.tasks import roll_dice_async

# Flask アプリの作成
app = Flask(__name__)

# Celery 設定
app.config["CELERY_BROKER_URL"] = "redis://localhost:6379/0"
app.config["CELERY_RESULT_BACKEND"] = "redis://localhost:6379/0"
celery = Celery(app.name, broker=app.config["CELERY_BROKER_URL"])
celery.conf.update(app.config)

# OpenTelemetry の設定
tracer = trace.get_tracer("diceroller.tracer")


@app.route("/rolldice")
def roll_dice():
    username = get_username()

    # トレースコンテキストを取得して Celery に渡す
    headers = {}
    inject(headers)
    roll_dice_async.apply_async(kwargs={"username": username, "headers": headers})

    return "Rolling dice asynchronously!"


def get_username():
    with tracer.start_as_current_span("get_username") as span:
        username = "hawksnowlog"
        span.set_attribute("username", username)
        return username


if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8080)

lib/tasks.py

from random import randint

from opentelemetry import trace
from opentelemetry.propagate import extract

from lib.worker import celery

tracer = trace.get_tracer("diceroller.tracer")


@celery.task
def roll_dice_async(username, headers):
    return roll(username, headers)


def roll(username: str, headers):
    # OpenTelemetry のコンテキストを受け取り、トレースを継続
    with tracer.start_as_current_span("roll", context=extract(headers)) as span:
        res = randint(1, 6)
        span.set_attribute("roll.value", res)
        return res

lib/worker.py

from celery import Celery
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.instrumentation.celery import CeleryInstrumentor
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor

# Celery インスタンスを作成
celery = Celery("tasks", broker="redis://localhost:6379/0")
celery.conf.update(broker_connection_retry_on_startup=True)

# OpenTelemetry の設定
tracer_provider = TracerProvider()
tracer_provider.add_span_processor(
    BatchSpanProcessor(OTLPSpanExporter(endpoint="http://localhost:4317"))
)

# Celery の OpenTelemetry を有効化
CeleryInstrumentor().instrument()

起動

Flask

  • export OTEL_PYTHON_LOGGING_AUTO_INSTRUMENTATION_ENABLED=true && pipenv run opentelemetry-instrument --logs_exporter otlp --service_name dice-server python app.py

Celery

  • export OTEL_PYTHON_LOGGING_AUTO_INSTRUMENTATION_ENABLED=true && pipenv run opentelemetry-instrument --logs_exporter otlp --service_name dice-server celery -A lib.tasks worker --loglevel=INFO

Jaeger

  • docker run --rm --name jaeger -e COLLECTOR_ZIPKIN_HOST_PORT=:9411 -p 16686:16686 -p 4317:4317 -p 4318:4318 -p 14250:14250 -p 14268:14268 -p 14269:14269 -p 9411:9411 jaegertracing/all-in-one:latest

動作確認

localhost:8080/rolldice にアクセスして Jaeger にアクセスすると一つのトレース内に Flask と Celery のタスクが処理が含まれていることが確認できると思います

最後に

Flask + Celery に OLTP のトレースを設定する方法を紹介しました

トレースに間があるのは非同期特有なのかもしれません
もしくは BatchSpanProcessor を使っているのでリアルタイムでなくバッチ処理として送信している影響かもしれません

参考サイト

2024年1月19日金曜日

Django と celery の連携方法

Django と celery の連携方法

概要

Django と Celery を連携してみました
公式のドキュメントだと django 側には celery 連携のドキュメントなく celery 側にあるようです

環境

  • macOS 11.7.10
  • Python 3.11.6
  • Django 5.0.1
  • celery 5.3.6
  • Redis 7.2.1

mysite/settings.py

一部抜粋

# for celery
CELERY_CACHE_BACKEND = "default"
CELERY_BROKER_URL = "redis://localhost:6379"

polls/celery.py

import os

from celery import Celery

os.environ.setdefault("DJANGO_SETTINGS_MODULE", "mysite.settings")

app = Celery("polls")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks()


@app.task(bind=True, ignore_result=True)
def debug_task(self):
    print(f"Request: {self.request}")

polls/__init__.py

from .celery import app as celery_app

__all__ = ("celery_app",)

ワーカー起動

  • pipenv run celery -A polls worker -l INFO

動作確認

  • pipenv run python manage.py shell
from polls.celery import debug_task
debug_task.apply_async()

最後に

連携自体はそれほど難しくないようです
設定やタスクは celery の公式にあるように記載しましたが別ファイルで管理してもいいかなと重ます

参考サイト

2023年3月20日月曜日

(celery) flower にベーシック認証をつける方法

(celery) flower にベーシック認証をつける方法

概要

flower は celery の管理ツールです
今回は flower にベーシック認証を付与する方法を紹介します
なお今回は docker で起動します

環境

  • Ubuntu 18.04
  • docker 20.10.7
  • Python 3.10.2
  • flower 1.2

ベーシック認証用の引数を付与して起動

--basic_auth オプションを使うだけです

  • docker run -p 5555:5555 --rm mher/flower celery --broker=redis://172.17.0.1:6379 flower --address=0.0.0.0 --pt=5555 --basic_auth=user1:password1

動作確認

これで localhost:5555 にアクセスするとベーシック認証が発動するのが確認できると思います

最後に

引数で簡単に指定できるのはいいのですが平文なので何ともという感じです

もしこれが嫌な場合は nginx などのリバースプロキシ配下で flower を起動する方法があるのでそれを使う感じになるかなと思います

参考サイト

2023年1月27日金曜日

Celery の visibility timeout の挙動を確認してみた

Celery の visibility timeout の挙動を確認してみた

概要

タイトルの通りです

visibility timeout に設定した秒数が経過した場合にタスクがどうなるかの影響を試してみました

環境

  • macOS 11.7.2
  • Python 3.10.2
    • celery v5.2.7

サンプルコード

result_backend = "redis://localhost"
broker_url = "redis://localhost"
worker_prefetch_multiplier = 1
task_acks_late = True
result_backend_transport_options = broker_transport_options = { 'visibility_timeout': 5 }
import time
from celery import Celery

app = Celery("tasks")
app.config_from_object("celeryconfig")


@app.task(bind=True)
def add(self, x, y):
    return x + y


@app.task(bind=True)
def multi(self, x, y):
    time.sleep(1800)
    return x * y

起動

  • pipenv run celery -A tasks worker --loglevel=info
  • vim app.py

from celery import chain
from tasks import add, multi

tasks = chain(add.s(1, 2),
              multi.s(10),
              multi.s(2)).apply_async()

print(tasks.get())

起動

動作確認

実行時のログは以下のようになりました

[2023-01-27 11:50:19,854: INFO/MainProcess] Connected to redis://localhost:6379//
[2023-01-27 11:50:19,862: INFO/MainProcess] mingle: searching for neighbors
[2023-01-27 11:50:20,875: INFO/MainProcess] mingle: all alone
[2023-01-27 11:50:20,894: INFO/MainProcess] celery@node1.local ready.
[2023-01-27 11:50:24,696: INFO/MainProcess] Task tasks.add[f4ca797b-1a76-4800-8a5f-ce7e39ad727d] received
[2023-01-27 11:50:24,718: INFO/MainProcess] Task tasks.multi[322d8025-e7ea-43d6-8f7f-e54008d46db4] received
[2023-01-27 11:50:24,719: INFO/ForkPoolWorker-2] Task tasks.add[f4ca797b-1a76-4800-8a5f-ce7e39ad727d] succeeded in 0.019410155015066266s: 3
[2023-01-27 11:51:49,898: INFO/MainProcess] Task tasks.multi[322d8025-e7ea-43d6-8f7f-e54008d46db4] received
[2023-01-27 11:53:29,947: INFO/MainProcess] Task tasks.multi[322d8025-e7ea-43d6-8f7f-e54008d46db4] received

ちょっと解説

どうやら visibiity_timeout に設定した時間が経過してすぐに再キューされるわけではないようです

何度か試してみましたがまちまちで visibility_timeout = 5 のときは 1 分から 2 分後に再キューされるような挙動でした
なので visibility_timeout = 5 だからと言って time.speep(10) くらいにしてしまうと再キューが発生することなく正常に終了してしまいます

また ack_late = True のオプション (もしくは task_acks_late = True) に設定していないと visibility_timeout の時間が反映されないようなのでそこも注意が必要です

最後に

プリフェッチを無効にする場合に task_acks_late = True の設定が必須なのでその場合には visibiliti_timeout の設定も合わせて調整する必要がありそうです

参考サイト

2023年1月23日月曜日

Celery カスタムクラスとして作成したタスクでハンドラの挙動を確認した

Celery カスタムクラスとして作成したタスクでハンドラの挙動を確認した

概要

前回 Celery のカスタムタスクを作成する方法を紹介しました
今回はカスタムタスクを使って各種イベントハンドラを使用する方法を紹介します
エラーやリトライ時の挙動をハンドリングできます

環境

  • macOS 11.7.2
  • Python 3.10.2
  • celery 5.2.7

サンプルコード

  • vim tasks.py
from celery import Celery

app = Celery("tasks")
app.config_from_object("celeryconfig")


class AddTask(app.Task):

    def __init__(self):
        self.name = "AddTask"

    def run(self, x, y, *args, **kwargs):
        # raise RuntimeError  # call to on_failure
        # self.retry()  # call to on_retry
        return x + y

    def before_start(self, task_id, *args, **kwargs):
        print("before_start")
        print(task_id)

    def after_return(self, status, retval, task_id, einfo = None, *args, **kwargs):
        print("after_return")
        print(status)

    def on_failure(self, exc, task_id, einfo = None, *args, **kwargs):
        print("on_failure")
        print(exc)

    def on_retry(self, exc, task_id, einfo = None, *args, **kwargs):
        print("on_retry")
        print(exc)

    def on_success(self, retval, task_id, *args, **kwargs):
        print("on_success")
        print(retval)


class MultiTask(app.Task):

    def __init__(self):
        self.name = "MultiTask"

    def run(self, x, y, *args, **kwargs):
        return x * y


app.register_task(AddTask())
app.register_task(MultiTask())
  • vim celeryconfig.py
result_backend = "redis://localhost"
broker_url = "redis://localhost"
worker_prefetch_multiplier = 1
task_acks_late = True

実行コード

  • vim app.py
from celery import chain
from tasks import AddTask, MultiTask


class Base():

    def __init__(self):
        first = AddTask().s(1, 2)
        second = MultiTask().s(10)
        third = MultiTask().s(2)
        self.tasks = chain(first,
                           second,
                           third)

    def run(self):
        return self.tasks.apply_async()

    def append_task(self, task):
        self.tasks.tasks.append(task)


class Workflow1(Base):

    def __init__(self):
        forth = MultiTask().s(3)
        super().__init__()
        self.append_task(forth)


if __name__ == '__main__':
    base = Base()
    print(base.run().get())
    wf1 = Workflow1()
    print(wf1.run().get())

ポイント

before_start と after_return はタスクが成功しようが失敗しようが必ずコールされます
他の on_success, on_failure, on_retry はその間で呼ばれるハンドラになります

リトライはデフォルトだと3回行います
失敗したあとに180秒待ったあとリトライを行います
on_retry が 3 回呼び出されて MaxRetriesExceededError になると最後に on_failure が呼び出されます

参考サイト

2023年1月21日土曜日

Celery5 でタスクをクラスとして定義する方法

Celery5 でタスクをクラスとして定義する方法

概要

過去 に紹介した方法と少し違った方法を紹介します

環境

  • macOS 11.7.2
  • Python 3.10.2
  • celery 5.2.7

サンプルコード

  • vim tasks.py
from celery import Celery

app = Celery("tasks")
app.config_from_object("celeryconfig")


class AddTask(app.Task):

    def __init__(self):
        self.name = "AddTask"

    def run(self, x, y, *args, **kwargs):
        return x + y


class MultiTask(app.Task):

    def __init__(self):
        self.name = "MultiTask"

    def run(self, x, y, *args, **kwargs):
        return x * y


app.register_task(AddTask())
app.register_task(MultiTask())
  • vim celeryconfig.py
result_backend = "redis://localhost"
broker_url = "redis://localhost"
worker_prefetch_multiplier = 1
task_acks_late = True

実行コード

  • vim app.py
from celery import chain
from tasks import AddTask, MultiTask


class Base():

    def __init__(self):
        first = AddTask().s(1, 2)
        second = MultiTask().s(10)
        third = MultiTask().s(2)
        self.tasks = chain(first,
                           second,
                           third)

    def run(self):
        return self.tasks.apply_async()

    def append_task(self, task):
        self.tasks.tasks.append(task)


class Workflow1(Base):

    def __init__(self):
        forth = MultiTask().s(3)
        super().__init__()
        self.append_task(forth)


if __name__ == '__main__':
    base = Base()
    print(base.run().get())
    wf1 = Workflow1()
    print(wf1.run().get())

少し解説

Celery から作成した app オブジェクトの Task クラスを継承して作成します

name フィールドは必須です
run メソッドも必須です

定義したタスククラスは 実行する場合は app.register_task を使って登録します
この登録方法が前回少し異なっていました
app.tasks.register を使うとなぜか AttributeError: 'NoneType' object has no attribute 'push' というエラーが発生して request_stack.push ができなくタスクが登録できませんでした

タスクを呼び出す場合は普通にクラスを生成する感じで呼び出します
あとはこれまで通りと同じように s() や chain を使ってタスクを扱います

動作確認

  • pipenv run celery -A tasks worker --loglevel=info
  • pipenv run python app.py

最後に

Celery のタスクのクラス化はドキュメントが少ない印象があるのでちょくちょく紹介していこうかなと思っています

2023年1月11日水曜日

celery のタスクで chain 作成後にタスクを挿入する方法

celery のタスクで chain 作成後にタスクを挿入する方法

概要

chain は celery でタスクを順番に実行することができるワークフロー機能です
chain で作成したタスクの一覧に対して実行する前にタスクの順番を変えたい場合があると思います
そんなときのテクニックを紹介します

環境

  • macOS 11.7.2
  • Python 3.10.2
  • celery 5.2.7

タスク

from celery import Celery

app = Celery("tasks")
app.config_from_object("celeryconfig")


@app.task(bind=True)
def add(self, x, y):
    return x + y


@app.task(bind=True)
def multi(self, x, y):
    return x * y

サンプルコード

from celery import chain
from tasks import add, multi

first = add.s(1, 2)
second = multi.s(10)
third = multi.s(2)
forth = multi.s(3)

tasks = chain(first,
              second,
              third)
# tasks.tasks.insert(2, forth)
tasks.tasks.append(forth)

result = tasks.apply_async()
print(result.get())

ちょっと解説

chain で生成されたオブジェクトは _chain クラスのオブジェクトになります
このクラスにある tasks プロパティは配列で管理されておりこの中にタスクの一覧が入っています
あとはこの配列に対してタスクの追加などを行うだけです

もう少し汎用的にしてみる

この仕組みを応用してもう少し汎用的な形にしてみます
基本となる chain のタスク一覧を管理するクラスを作成しそのタスク一覧をもとに別の chain を実行するワークフローを作成します

from celery import chain
from tasks import add, multi


class Base():

    def __init__(self):
        first = add.s(1, 2)
        second = multi.s(10)
        third = multi.s(2)
        self.tasks = chain(first,
                           second,
                           third)

    def run(self):
        return self.tasks.apply_async()

    def append_task(self, task):
        self.tasks.tasks.append(task)


class Workflow1(Base):

    def __init__(self):
        forth = multi.s(3)
        super().__init__()
        self.append_task(forth)


if __name__ == '__main__':
    base = Base()
    print(base.run().get())
    wf1 = Workflow1()
    print(wf1.run().get())

こんな感じに記載することで chain で実行するタスクのリストを重複することなく管理する ことができるようになります

最後に

さらにタスクの一覧と chain オブジェクトを別のクラスとして管理してもいいかなと思います

テストも書きやすくなるかなと思います

2022年11月15日火曜日

CeleryでRetryしたジョブを見つける方法

CeleryでRetryしたジョブを見つける方法

概要

Celery はリトライしたジョブは最終的に SUCCESS or FAILURE になるためリトライした RETRY のステータスのジョブは残りません

成功したがどのジョブがリトライしたかを調べたい場合には flower の API を使うと簡単です

環境

  • Ubuntu 18.04
  • Python 3.10.2
  • flower 1.2

サンプルコード

import requests

api_root = 'http://localhost:5555/api'
task_api = '{}/tasks'.format(api_root)

res = requests.get(task_api)
tasks = res.json()
for task_id, info in tasks.items():
    if info['retries'] > 0:
        print(task_id)
        print(info)

参考サイト

2022年10月7日金曜日

Celeryでプリフェッチ(RECEIVED)されないための設定

Celeryでプリフェッチ(RECEIVED)されないための設定

概要

celery にはプリフェッチという機能がありワーカーが実行中であっても次のジョブを事前に取得して RECEIVED として管理する機能があります
RECEIVED になったジョブはワーカーが保持するすべてのジョブが終了するまで実行されないので長い処理を実行中のワーカーにプリフェッチされると RECEIVED で永遠に実行されないジョブが発生してしまいます
また RECEIVED の状態で1時間経過すると再度ジョブをエンキューするという謎の仕様がありこのせいで同一ジョブが複数回実行されるというバグもあるようです (参考)

今回のこのプリフェッチを無効にする方法があるとのことなので試してみました

環境

  • macOS 11.7
  • Python 3.10.2
  • celery 5.2.7
  • flower 1.2.0

設定ファイル

worker_prefetch_multipliertask_acks_late を設定します

  • vim celeryconfig.py
result_backend = 'redis://localhost'
broker_url = 'redis://localhost'
worker_prefetch_multiplier = 1
task_acks_late = True

タスク

プリフェッチされているかを確認するためにあえて wait を入れています

  • vim tasks.py
from celery import Celery

app = Celery('tasks')
app.config_from_object('celeryconfig')

@app.task
def add(x, y):
    import time
    time.sleep(3600)
    return x + y

@app.task
def multi(x, y):
    return x * y

メイン

  • app.py
from celery import chain
from tasks import add, multi

tasks = chain(add.s(1, 2),
              multi.s(10),
              multi.s(2)).apply_async()

動作確認

  • pipenv run celery -A tasks worker --loglevel=info

ワーカーを起動しいくつかのジョブを登録します

  • for i in `seq 0 5`; do pipenv run python app.py; done

すると4つのジョブはSTARTEDになりますが残り2つはプリフェッチされず RECEIVED にもならないことが確認できます

ちょっと流れを解説

redis 内で管理されているキーの説明を少しします

STARTED になったタスクは unacked と unacked_index というキーで管理されています
それぞれ Hash 型と ZSet 型で管理されています

まだどのワーカーにも属していないタスクは celery というキー (キュー) 内で管理されています
celery は List 側になっています

最後に

注意点としては worker_prefetch_multiplier が 3,4 系では使えない点です (4.2 から使えるようになった?)

詳細は不明ですが素直に5系にアップグレードしてから使ったほうが良いかなと思います

プリフェッチすらされていないジョブのタイムアウトなどあるのか気になりました

参考サイト

2022年10月5日水曜日

Celery の状況を UI で管理できる flower を試してみた

Celery の状況を UI で管理できる flower を試してみた

概要

celery には flower という管理 UI の機能があります
今回はインストールから実際に起動して動作確認するところまで行ってみました

環境

  • macOS 10.15.5
  • Python 3.8.3
    • celery 4.4.4
    • flower 0.9.4 (1.2.0)

flower のインストール

  • pipenv install flower

flower の起動

  • pipenv run flower -A tasks --broker=redis://localhost:6379 --port=5555

--broker で celery が使っているブローカを指定します
--port でポートを指定できます
デフォルトは 5555 になっています

P.S 20221005

flower1.2.0では以下のように起動コマンドが変わっているのでご注意ください

  • pipenv run celery -A tasks flower --broker=redis://localhost:6379 --port=5555

動作確認

今回使うタスクは以下の通りです

from celery import Celery

app = Celery('tasks', broker='redis://localhost:6379')

@app.task
def add(x, y):
    print(x + y)

起動しましょう

  • pipenv run celery -A tasks worker -l info

あとは何回かタスクを消化させてみます

from tasks import add

add.delay(100, 1)
  • pipenv run python test.py

これで localhost:5555 にアクセスしてみましょう
起動しているワーカの一覧が表示されると思います

ワーカをクリックすると詳細情報が確認できます
プールサイズの変更なども画面からできるようです
またタスクを実行するキュー (Consumer) も画面から追加できるようです

最後に

celery の管理 UI である flower を試してみました
ブローカを指定して起動するだけなので簡単です
ある程度のワーカの設定変更もできるのでプロセスを再起動することなく操作できるのは嬉しい点かなと思います

認証を追加することもできるようなので特定のユーザにのみ閲覧させることも可能です

参考サイト

2022年6月9日木曜日

celery の chain でタスクごとにキューをルーティングする方法

celery の chain でタスクごとにキューをルーティングする方法

概要

過去にキューをルーティングする方法を紹介しました
今回は chain と組み合わせてルーティングする方法を紹介します

環境

  • macOS 11.6.6
  • Python 3.10.2
  • celery 5.2.7

タスク

まずは普通にタスクを定義します
ここではキューは指定しません

  • vim tasks.py
from celery import Celery

app = Celery('sub_tasks', backend='redis://localhost', broker='redis://localhost')

@app.task
def add(x, y):
    return x + y

@app.task
def multi(x, y):
    return x * y

chain

次に chain を使ってタスク同士をつなげます
このときにタスクがエンキューされるキューを指定します
apply_async では指定せずにタスクに対して set をコールすることでタスクに紐付けるキューを指定できます

  • vim app.py
from celery import chain
from tasks import add, multi

tasks = chain(add.s(1, 2),
              multi.s(10).set(queue='multi'),
              multi.s(2)).apply_async()

print(tasks.get())

上記の場合は 1 つ目と 3 つ目のタスクはデフォルトのキューで 2 つ目のタスクは「multi」というキューを使います

ワーカーをそれぞれ起動

デフォルトのキューに入ったタスクを処理するワーカーと multi キューに入ったタスクを処理するワーカーを起動します

  • pipenv run celery -A tasks worker -l info
  • pipenv run celery -A tasks worker -Q multi -l info

動作確認

エンキューして動作確認します

  • pipenv run python app.py

これでキューを指定したそれぞれのワーカーでタスクが処理されていることが確認できると思います

apply_async 時に指定するキューの意味

デフォルトのキューになります
タスクそれぞれでキューを set していない場合に使用されるキューになります

2022年1月24日月曜日

Celery でタスクが登録されないときに確認すること (コードで説明編)

Celery でタスクが登録されないときに確認すること (コードで説明編)

概要

前回文章で説明しました

実際にサンプルのコードが合ったほうがイメージしやすいのでコードでも紹介します

環境

  • macOS 11.6.2
  • Python 3.10.1
  • celery 5.2.3

登録されないパターン

まずはタスクが登録されないパターンです
これで celery を起動しても作成した add, multi はタスクとして登録されません

もしこの状態でエンキューされると処理するタスクがありませんというエラーになります

  • vim tasks/celery.py
from celery import Celery

app = Celery('tasks', backend='redis://localhost', broker='redis://localhost')
  • vim tasks/calc.py
from .celery import app

@app.task
def add(x, y):
    return x + y

@app.task
def multi(x, y):
    return x * y
  • pipenv run celery -A tasks.celery.app worker --loglevel=INFO

登録されるように import する

ではどうすれば別ディレクトリに定義したタスクが登録されるかというと celery 実行時にタスクを import してあげます

またこの import の際には循環参照に気をつけなければいけません

一応これでも動きますがすぐに循環参照しそうなコードになっています

  • vim tasks/celery.py
from celery import Celery
# from tasks.calc import add, multi <- ここだと循環参照になる

app = Celery('tasks', backend='redis://localhost', broker='redis://localhost')

from tasks.calc import add, multi
  • vim tasks/calc.py
from .celery import app

@app.task
def add(x, y):
    return x + y

@app.task
def multi(x, y):
    return x * y
  • pipenv run celery -A tasks.celery.app worker --loglevel=INFO

タスクをロードするモジュールを別で作成するパターン

celery オブジェクト管理するモジュールとタスクをロードするモジュールを分けます
celery オブジェクトを複数のモジュールで持つことになりますがこれでも可能です

celery 実行時に loader にある celery オブジェクト (tasks.loader.app ) を指定します

こちらのパターンのほうが循環参照を避けられるかなと思います

  • vim tasks/celery.py
from celery import Celery

app = Celery('tasks', backend='redis://localhost', broker='redis://localhost')
  • vim tasks/loader.py
from tasks.calc import add, multi
from celery import Celery

app = Celery('tasks', backend='redis://localhost', broker='redis://localhost')
  • vim tasks/calc.py
from .celery import app

@app.task
def add(x, y):
    return x + y

@app.task
def multi(x, y):
    return x * y
  • pipenv run celery -A tasks.loader.app worker --loglevel=INFO

動作確認

上記すべて以下のスクリプトで動作確認できます
一番上のパターンはタスクが登録されていないのでエラーになります

  • vim main.py
from celery import chain
from tasks.calc import add, multi

tasks = chain(
    add.s(1, 2),
    multi.s(10)).apply_async()

print(tasks.get())

最後に

更にこれに flask などの Webフレームワークや SQLAlchemy などの ORM が加わってくると app コンテキストが必要になってきます

モジュールを分ける場合にはそれらが循環参照しないように気をつける必要があるので大変です

2022年1月21日金曜日

Celery でタスクが登録されないときに確認すること

Celery でタスクが登録されないときに確認すること

概要

celery を実行した際にタスクが登録されない場合にチェックする点を紹介します

環境

  • macOS 11.6.2
  • Python 3.10.1
  • celery 5.2.3

タスクがちゃんとどこかで import されているか

celery -A app.tasks.celery などで実行した場合に app.tasks や app.tasks が自動で読み込むモジュール (例えば app/__init__.py など) でちゃんとタスクとして定義したモジュールが import されている確認しましょう

実行したタスクモジュールを読み込んだ際にそのタスクどこかで import されいないと celery はタスクとして登録してくれません

もしくは app.tasks 内で定義した celery オブジェクトを使って @celery.task() アノテーションでタスク登録しても OK です

クラスの場合は register_task をする

クラスとしてタスクを定義した場合は Celery.register_task() で必ずクラスを登録しましょう

もしファイルを分ける場合はタスクはクラスとして定義したほうが管理の面ではいいと思います

2021年11月12日金曜日

celery の s() と si() の使い分け

celery の s() と si() の使い分け

概要

s() は前のタスクの結果を受け取ります
si() は前のタスクの結果を受け取りません

環境

  • macOS 11.6.1
  • Python 3.8.12
    • celery 5.1.2

タスクの定義

  • vim sub_tasks.py
from celery import Celery

app = Celery('sub_tasks', backend='redis://localhost', broker='redis://localhost')

@app.task
def add(x, y):
    return x + y

@app.task
def multi(x, y):
    return x * y

s() を使う場合

まずは s() を使います
前のタスクの結果は次のタスクの第一引数に渡されます

from celery import chain
from sub_tasks import add, multi

tasks = chain(
    add.s(1, 2),
    multi.s(10)).apply_async()

print(tasks.get())

結果は 1+2=3 -> 3 x 10 = 30 になります

以下のような使い方はエラーになります

from celery import chain
from sub_tasks import add, multi

tasks = chain(
    add.s(1, 2),
    multi.s(10, 10)).apply_async()

print(tasks.get())

TypeError: multi() takes 2 positional arguments but 3 were given 前のタスクの結果と合わせて引数が足りないと言われてエラーになります

si() を使う場合

si() を使うことで前のタスクの結果を使わないようにできます

from celery import chain
from sub_tasks import add, multi

tasks = chain(
    add.s(1, 2),
    multi.si(10, 10)).apply_async()

print(tasks.get())

結果は 10 x 10 = 100 になります
1+2=3 の結果は使っていないことがわかります

以下の使い方はエラーになります

from celery import chain
from sub_tasks import add, multi

tasks = chain(
    add.s(1, 2),
    multi.si(10)).apply_async()

print(tasks.get())

TypeError: multi() missing 1 required positional argument: 'y'

前のタスクが None を返す場合は

None でも関係なく s() を使った場合は次のタスクの第一引数に入るので注意しましょう

どちらを使うのがいいのか

個人的には以下のようなルールで使うのがいいと思っています

  • 定義したタスクは必ず単独で使う -> si() を使う
  • 定義したタスクは前の結果を使う可能性がある -> s() を使う

前者の場合タスクを定義する際の引数に前のタスクの結果を受け取るための引数がなくなるので必ず si() を使い続ける必要があります

後者は前のタスクの結果を受け取るための引数がありながらもケースによっては引数を使わないこともできます
ただ使わない引数を定義することになるので少しコードが紛らわしくなる可能性があります

タスクとして柔軟なのは間違いなく後者です
例えば将来的に定義したタスクの前にタスクが来てその結果を使う場合には後者で実装しなければなりません
前にタスクが来てもタスクの結果を使うことがなければ si() でも OK です

最後に

si() = s() + immutable=True なので素直に s() を使うのがいいのかもしれません

タスクを定義する場合も必ず前のタスクを受け取るようにするのがいいのかもしれません

参考サイト