You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python:使用Celery在多服务器上处理参数列表的问题

Hey there! Let's break down how to solve this problem given your supercomputing cluster constraints—Celery is totally up to the task, we just need to lean into distributed state management and careful task configuration.

Core Approach: Centralized Task Queue + Distributed Locking

The key here is to ensure all your cluster nodes share a single source of truth for task state, and that each input is locked to one worker at a time. Here's a step-by-step implementation:

1. Set Up a Shared Message Broker & Result Backend

First, you need a centralized broker (like Redis or RabbitMQ) that all cluster nodes can reach. This will be the single queue where all tasks live, and a result backend to track task statuses.

Configure your Celery app to use this shared resource:

from celery import Celery

# Replace with your shared broker/backend host (must be accessible to all cluster nodes)
app = Celery(
    'cluster_task_processor',
    broker='redis://shared-redis-node:6379/0',
    backend='redis://shared-redis-node:6379/0'
)

This ensures every worker node pulls tasks from the same global queue—no duplicate task distribution out of the gate.

2. Add Distributed Locking to Prevent Duplicate Processing

Even with a shared queue, edge cases (like a node being interrupted mid-task) can lead to potential duplicates. Use a Redis-backed lock to ensure each input is only processed by one worker at a time.

We'll use the celery-redis-lock library for simplicity:
First install it:

pip install celery-redis-lock

Then update your task to include the lock:

from celery_redis_lock import lock
from celery.utils.log import get_task_logger

logger = get_task_logger(__name__)

@app.task(bind=True, autoretry_for=(Exception,), retry_backoff=3, max_retries=5)
def process_single_input(self, input_item):
    # Use the input itself as a unique lock key (adjust if inputs aren't unique strings)
    lock_key = f"task_lock:{input_item}"
    
    # Lock for 2x your expected max task runtime to cover processing delays
    with lock(lock_key, timeout=7200):
        # Double-check if the task was already completed (in case lock expired)
        if self.backend.get_status(self.request.id) == 'SUCCESS':
            logger.info(f"Skipping already processed input: {input_item}")
            return f"Input {input_item} already processed"
        
        # Your actual input processing logic here
        logger.info(f"Processing input: {input_item}")
        processing_result = your_input_processor_function(input_item)
        
        # Mark task as completed (optional but helpful for status checks)
        self.update_state(state='SUCCESS', meta={'result': processing_result})
        return processing_result

The lock ensures only one worker can execute the task for a given input, and the retry logic handles cases where a node drops mid-processing.

3. Preload All Tasks Into the Queue Upfront

Since you can only send one startup command per server, populate the task queue once before any nodes are scheduled. Write a simple script to push all your inputs into the shared queue:

# populate_tasks.py
from cluster_task_processor import process_single_input

# Replace with your full list of inputs
INPUT_LIST = ["input_001", "input_002", "input_003", ...]

for input_item in INPUT_LIST:
    # Send each input as a separate task to the queue
    process_single_input.delay(input_item)

Run this script once (from any machine that can reach the broker) to load all tasks into the queue. Now, whenever a cluster node starts your worker, it will immediately pull tasks from this queue.

4. Configure Celery Workers for Cluster Constraints

When you send the startup command to each server, launch a Celery worker with settings tailored to your cluster's single-process limit:

celery -A cluster_task_processor worker --concurrency=1 --loglevel=info
  • --concurrency=1 ensures each worker runs only one task at a time (matches your one-command-per-server constraint)
  • Adjust --loglevel to warning or error if you don't need verbose logs

5. Verify Completion & Handle Edge Cases

To make sure no inputs are missed, write a quick status check script:

# check_completion.py
from cluster_task_processor import process_single_input
from celery.result import AsyncResult

INPUT_LIST = ["input_001", "input_002", ...]
unprocessed_inputs = []

# For large lists, tracking processed inputs in a Redis set is more efficient
# Example: if not redis_client.sismember("processed_inputs", input_item): unprocessed_inputs.append(input_item)

# For this example, we'll assume you have a mapping of inputs to task IDs
for input_item in INPUT_LIST:
    task_id = get_task_id_for_input(input_item)
    task_result = AsyncResult(task_id)
    if task_result.state != 'SUCCESS':
        unprocessed_inputs.append(input_item)

if unprocessed_inputs:
    print(f"Found unprocessed inputs: {unprocessed_inputs}")
    # Resend the unprocessed tasks to the queue
    for item in unprocessed_inputs:
        process_single_input.delay(item)
else:
    print("All inputs processed successfully!")

内容的提问来源于stack exchange,提问作者duhaime

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 04:07:17