0Pricing
AI Agents · レッスン

非同期ワークフローとバックグラウンドジョブ

長時間実行するエージェントにはキュー(Celery、RQ、Temporal)が必要です。ジョブidを返し、ステータスをポーリングします。

「非同期ワークフローとバックグラウンドジョブ」はCoddyKit上の無料AI Agentsレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Agents学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Agentsコースには全4レッスンが含まれています。

このレッスンの一部はまだ翻訳されておらず、英語で表示されています。

リクエストサイクルが遅すぎる場合

エージェントのタスクには、30秒、5分、あるいは数時間かかるものがあります。HTTP接続を開いたままにはできません。非同期のバックグラウンドジョブに移してください。

パターン

  1. クライアントがタスクをPOST -> サーバーがジョブを作成し、job_idを返す
  2. ワーカーがキューからジョブを取得する
  3. クライアントがGET /jobs/{id}をポーリングしてステータスを確認する
  4. 完了すると、サーバーが結果を返す

Submit Endpoint

from uuid import uuid4

@app.post('/jobs')
def submit(req: AgentRequest):
    job_id = str(uuid4())
    redis.hset(f'job:{job_id}', mapping={'status': 'queued', 'user_id': req.user_id})
    queue.enqueue('run_agent_job', job_id, req.query)
    return {'job_id': job_id, 'status': 'queued'}

Status Endpoint

@app.get('/jobs/{job_id}')
def status(job_id: str):
    data = redis.hgetall(f'job:{job_id}')
    if not data:
        raise HTTPException(404)
    return {
        'status': data['status'],
        'result': data.get('result'),
        'error': data.get('error')
    }

ワーカー

別のプロセスがキューからジョブを取得します。

def run_agent_job(job_id, query):
    redis.hset(f'job:{job_id}', 'status', 'running')
    try:
        result = run_agent(query)
        redis.hset(f'job:{job_id}', mapping={'status': 'done', 'result': result})
    except Exception as e:
        redis.hset(f'job:{job_id}', mapping={'status': 'failed', 'error': str(e)})

キューの選択肢

  • RQ(Redis Queue)— 最小構成のPython向け
  • Celery — Pythonの定番で、多くの機能を備える
  • Temporal — 永続的なワークフロー、リトライ、可観測性
  • Dramatiq — Celeryの代替
  • クラウドネイティブ — Cloud Tasks、SQS、Pub/Sub

永続的なワークフローのためのTemporal

エージェントのワークフローは状態を持ちます。Temporalの永続的な実行モデルは、こうした用途に非常に適しています。

import temporalio

@temporalio.workflow.defn
class AgentWorkflow:
    @temporalio.workflow.run
    async def run(self, query: str) -> str:
        plan = await workflow.execute_activity(plan_step, query, schedule_to_close_timeout=timedelta(minutes=2))
        results = await workflow.execute_activity(execute_step, plan, schedule_to_close_timeout=timedelta(minutes=10))
        return await workflow.execute_activity(synthesise_step, results)

部分更新のストリーミング

長時間実行されるジョブでは、進捗の更新が役立ちます。Server-Sent EventsまたはWebSocketsを使用してください。

@app.get('/jobs/{job_id}/stream')
def stream_progress(job_id):
    def gen():
        while True:
            update = redis.brpop(f'updates:{job_id}', timeout=30)
            if not update:
                yield 'data: {"status": "timeout"}\n\n'
                break
            yield f'data: {update[1].decode()}\n\n'
            if 'done' in update[1].decode():
                break
    return StreamingResponse(gen(), media_type='text/event-stream')

ジョブのTTL

ジョブの記録を永遠に保持しないでください。

redis.expire(f'job:{job_id}', 86400)   # 1 day

リトライ

一時的な障害は自動的にリトライし、恒久的な障害は人が確認するためにDLQへ送ります。

import time

def retry(retries=3, retry_backoff=True):
    def decorator(func):
        def wrapper(*args, **kwargs):
            for attempt in range(1, retries + 1):
                try:
                    return func(*args, **kwargs)
                except Exception as e:
                    print(f'attempt {attempt} failed: {e}')
                    if attempt == retries:
                        raise
            return None
        return wrapper
    return decorator

attempts = {'n': 0}

@retry(retries=3, retry_backoff=True)
def run_agent_job(job_id):
    attempts['n'] += 1
    if attempts['n'] < 3:
        raise RuntimeError('transient error')
    return f'job {job_id} done'

print(run_agent_job('job-1'))

同時実行数の制限

GPUやAPIのコストを管理するため、同時に実行するワーカー数を制限してください。

rq worker --burst --max-jobs 1000 --queue agent

# Or per-queue concurrency in Temporal worker config.

ユーザーごとのクォータ

ユーザーごとに実行中のジョブ数を追跡してください。

key = f'inflight:{user_id}'
if redis.scard(key) >= 5:
    raise HTTPException(429, 'Too many in-flight jobs')
redis.sadd(key, job_id)

可観測性

各ジョブを最初から最後までトレースしてください。job_idとuser_idをスパンのタグに付けます。失敗したジョブでは、チケットまたはアラートを自動的に作成するようにしてください。

ステータスポーリングパターン

長時間実行されるエージェントに、送信とステータスポーリングのパターンを使うのはなぜですか。

まとめ

ジョブをPOSTし、IDを返し、ワーカーがキューから処理し、クライアントがポーリングします。単純な処理にはRQ/Celery、状態を持つワークフローにはTemporalを使用します。TTL、リトライ、クォータ、トレースも設定します。

よくある質問

「非同期ワークフローとバックグラウンドジョブ」レッスンは無料ですか?

はい。「非同期ワークフローとバックグラウンドジョブ」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Agentsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Agentsコースには全4レッスンが含まれています。

「非同期ワークフローとバックグラウンドジョブ」で何を学びますか?

長時間実行するエージェントにはキュー(Celery、RQ、Temporal)が必要です。ジョブidを返し、ステータスをポーリングします。 ブラウザで直接実行するハンズオンコードでAI Agentsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

AI Agentsを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのAI Agentsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。

「非同期ワークフローとバックグラウンドジョブ」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このAI Agentsレッスンでコードを書いて実行できますか?

はい。すべてのAI Agentsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. APIの背後でエージェントを提供する
  2. 非同期ワークフローとバックグラウンドジョブ
  3. レート制限とクォータ管理
  4. エージェントのBlue-GreenおよびCanaryデプロイ
← AI Agentsに戻る