Pythonにおける非同期タスクキュー:Celeryの活用法

Celeryの概要

CeleryはPythonで実装された分散型タスクキューであり、非同期処理を可能にするツールです。ウェブアプリケーションのバックグラウンド処理、データ処理、スクレイピングなど、時間のかかるタスクを即座に実行するのではなく、タスクキューに投入し、バックグラウンドで処理することで、システムの応答性を高めることができます。

Celeryの主要コンポーネント

  • ブローカー (Broker): タスクを保管し、ワーカーに配信するためのメッセージングミドルウェアです。RabbitMQやRedisが一般的に利用されます。プロデューサー(タスクを生成する側)がタスクをキューに投入し、ブローカーがそれを管理します。
  • ワーカー (Worker): ブローカーからタスクを受け取り、実際に処理を実行するプロセス群です。キューを監視し、新しいタスクが到着すると取り出して実行します。
  • バックエンド (Backend): タスクの実行結果を保存するためのストレージです。データベース(PostgreSQL, MySQLなど)、Redis、RabbitMQ (RPC) などが利用できます。

Celeryのセットアップと利用

ここでは、RabbitMQをブローカーとして利用したCeleryの基本的なセットアップとタスク実行手順を解説します。

プロジェクト構成

以下のディレクトリ構造でプロジェクトを作成します。

my_celery_project/
├── venv/                      # 仮想環境
├── celery_config.py           # Celeryアプリケーション設定
├── task_definitions.py        # タスク定義ファイル
└── client_scripts/            # クライアントスクリプト群
    ├── client_add_task.py
    └── client_mul_task.py

1. 仮想環境の準備とパッケージインストール

Pythonの仮想環境を作成し、必要なパッケージをインストールします。

python3 -m venv venv
source venv/bin/activate
pip install celery==5.3.6  # Celeryのバージョンは適宜調整
pip install flower         # モニタリングツール

2. RabbitMQのセットアップ

CeleryのブローカーとしてRabbitMQをDockerで起動します。ここでは、`celery_rabbitmq`という名前でコンテナを起動し、ホスト名を`rabbit`と設定します。

docker run -d -p 5672:5672 -h rabbit --name celery_rabbitmq rabbitmq

次に、CeleryがRabbitMQに接続するためのユーザーと仮想ホストを作成し、必要な権限を設定します。ここでは、`rabbit_user`、`rabbit_pass`、`rabbit_vhost`を使用します。

docker exec -it celery_rabbitmq rabbitmqctl add_user rabbit_user rabbit_pass
docker exec -it celery_rabbitmq rabbitmqctl add_vhost rabbit_vhost
docker exec -it celery_rabbitmq rabbitmqctl set_user_tags rabbit_user celery
docker exec -it celery_rabbitmq rabbitmqctl set_permissions -p rabbit_vhost rabbit_user ".*" ".*" ".*"

3. Celeryアプリケーションの定義

Celeryアプリケーションを初期化し、ブローカーとバックエンドの設定を行います。celery_config.pyを作成します。

# my_celery_project/celery_config.py
from celery import Celery

# RabbitMQブローカーへの接続URL
# RabbitMQがDockerコンテナで動作し、ポートがホストにマッピングされている場合、'localhost'を使用。
BROKER_URL = "amqp://rabbit_user:rabbit_pass@localhost:5672/rabbit_vhost"

# タスク結果の保存先 (RPCバックエンドを使用)
RESULT_BACKEND = "rpc://"

# Celeryアプリケーションの初期化
app = Celery(
    "my_celery_project", # プロジェクトのモジュール名 (ディレクトリ名と一致させる)
    broker=BROKER_URL,
    backend=RESULT_BACKEND,
    include=["my_celery_project.task_definitions"] # タスク定義モジュールの指定
)

# オプション設定 (例: タイムゾーン、タスク開始の追跡)
app.conf.update(
    task_track_started=True,
    timezone='Asia/Tokyo',
    enable_utc=True,
)

4. 非同期タスクの定義

ワーカーによって実行される非同期タスクを定義します。task_definitions.pyを作成し、複数のタスクをここに記述します。

# my_celery_project/task_definitions.py
from .celery_config import app
import time
import random

@app.task
def add_two_numbers(val1: int, val2: int) -> int:
    """
    2つの数値を加算し、指定された遅延後に結果を返すタスク。
    """
    task_name = add_two_numbers.name
    print(f"[{task_name}] タスク開始: {val1} と {val2} の加算")
    
    # 処理時間をシミュレート (5秒から10秒の間でランダムに遅延)
    delay_sec = random.randint(5, 10)
    time.sleep(delay_sec)
    
    result = val1 + val2
    print(f"[{task_name}] タスク終了. 結果: {result}, 処理時間: {delay_sec}秒")
    return result

@app.task
def multiply_numbers_in_list(numbers: list) -> int:
    """
    数値リストの全ての要素を乗算し、指定された遅延後に結果を返すタスク。
    """
    task_name = multiply_numbers_in_list.name
    print(f"[{task_name}] タスク開始: リスト {numbers} の乗算")
    
    # 処理時間をシミュレート (10秒から20秒の間でランダムに遅延)
    delay_sec = random.randint(10, 20)
    time.sleep(delay_sec)
    
    product = 1
    for num in numbers:
        product *= num
    
    print(f"[{task_name}] タスク終了. 結果: {product}, 処理時間: {delay_sec}秒")
    return product

5. タスクのディスパッチ(実行依頼)

タスク定義ファイルをインポートし、.delay() メソッドを使用してタスクをCeleryキューにディスパッチします。client_scriptsディレクトリに以下のファイルを作成します。

加算タスクをディスパッチ

# my_celery_project/client_scripts/client_add_task.py
from my_celery_project.task_definitions import add_two_numbers
import time

print("--- 加算タスクをCeleryにディスパッチ中 ---")
# delay()でタスクを非同期実行
async_result = add_two_numbers.delay(123, 456)
print(f"ディスパッチされたタスクID: {async_result.id}")

# タスクの完了状態と結果を初回チェック
print(f"タスク完了状態 (初回): {async_result.ready()}")
print(f"現在の結果 (初回、未完了ならNone): {async_result.result}")

# タスクの完了を待機するための時間(加算タスクの最大処理時間より少し長めに設定)
print("タスク完了まで約12秒待機...")
time.sleep(12)

# タスクの完了状態と結果を再チェック
print(f"タスク完了状態 (再確認): {async_result.ready()}")
# .result プロパティはタスク完了までブロッキングする場合があるため、ready()で確認後にアクセス推奨
if async_result.ready():
    print(f"最終結果: {async_result.result}")
else:
    print("タスクはまだ完了していません。")

乗算タスクをディスパッチ

# my_celery_project/client_scripts/client_mul_task.py
from my_celery_project.task_definitions import multiply_numbers_in_list
import time

print("--- 乗算タスクをCeleryにディスパッチ中 ---")
# delay()でタスクを非同期実行
async_result = multiply_numbers_in_list.delay([2, 3, 4, 5])
print(f"ディスパッチされたタスクID: {async_result.id}")

# タスクの完了状態と結果を初回チェック
print(f"タスク完了状態 (初回): {async_result.ready()}")
print(f"現在の結果 (初回、未完了ならNone): {async_result.result}")

# タスクの完了を待機するための時間(乗算タスクの最大処理時間より少し長めに設定)
print("タスク完了まで約25秒待機...")
time.sleep(25)

# タスクの完了状態と結果を再チェック
print(f"タスク完了状態 (再確認): {async_result.ready()}")
if async_result.ready():
    print(f"最終結果: {async_result.result}")
else:
    print("タスクはまだ完了していません。")

Celery Workerの起動とタスク実行

1. Celery Workerの起動

my_celery_projectの親ディレクトリに移動し、Celery Workerを起動します。-Aオプションでアプリケーションモジュールを指定し、--concurrencyで同時実行ワーカー数を設定します。

celery -A my_celery_project worker --loglevel=info --concurrency=4

Workerが正常にRabbitMQに接続されると、ターミナルに以下のような情報が表示されます。

...
[INFO/MainProcess] connected to amqp://rabbit_user:**@localhost:5672/rabbit_vhost
[INFO/MainProcess] mingle: searching for neighbours
[INFO/MainProcess] mingle: all alone
[INFO/MainProcess] celery@yourhostname ready.
...

2. タスクの実行

Workerを起動したまま、別のターミナルを開きます。my_celery_projectの親ディレクトリに移動し、各クライアントスクリプトを実行してタスクをディスパッチします。

加算タスクの実行

python -m my_celery_project.client_scripts.client_add_task

乗算タスクの実行

python -m my_celery_project.client_scripts.client_mul_task

タスクがディスパッチされると、Workerのターミナルにタスクの受信と実行に関するログが表示され、指定された遅延時間後に結果が返されます。クライアントスクリプトの出力とWorkerのログを比較することで、非同期処理の流れを確認できます。

Flowerによるタスク監視

1. Flowerの起動

Celery Flowerは、Celeryのリアルタイム監視ツールです。my_celery_projectの親ディレクトリで以下のコマンドを実行してFlowerを起動します。

celery -A my_celery_project flower

2. Webインターフェースでの監視

ウェブブラウザで http://localhost:5555 にアクセスすると、Flowerのダッシュボードが表示されます。ここでは、実行中のタスク、完了したタスク、ワーカーの状態などをリアルタイムで確認できます。

  • Tasks (タスク): 各タスクのID、状態、引数、実行時間、結果などを詳細に確認できます。
  • Dashboard (ダッシュボード): ワーカーノードの数、実行されたタスクの総数、成功/失敗したタスクの割合など、システムの概要を一目で把握できます。

Flowerを活用することで、Celeryタスクのデバッグや運用監視が格段に容易になります。

タグ: Python Celery RabbitMQ 非同期処理 タスクキュー

9月12日 16:19 投稿