Background Job Processing with Celery and Database Partitioning Strategies

Execution Context and Import Rules

When executing Celery workers or scripts, the currrent working directory modifies Python's module search path. Relative imports may fail if the file served as the entry point is not imported within a package context.

Initiating Tasks Programmatically

Tasks can be dispatched from command-line scripts using apply_async for immediate queuing or eta for delayed execution.

# worker_config.py
from celery import Celery

app = Celery(
    'task_manager', 
    broker='redis://localhost:6379/0', 
    backend='redis://localhost:6379/1', 
    include=['jobs.tasks']
)
# jobs/tasks.py
from .worker_config import app

@app.task(bind=True)
def execute_logic(self, arg_a, arg_b):
    print(f'Computing {arg_a} + {arg_b}')
    return {'value': arg_a + arg_b}

@app.task
def reverse_data(x, y):
    return x - y

To trigger these operations manually:

# dispatch_handler.py
from jobs.tasks import execute_logic
from datetime import timedelta, datetime

# Immediate task queueing
execute_logic.delay(25, 15)

# Scheduled task dispatch
future = execute_logic.apply_async(
    args=[10, 20], 
    eta=datetime.utcnow() + timedelta(seconds=30)
)
print(f'Task Identifier: {future.id}')

Retrieving results requires the AsyncResult interface:

# fetch_status.py
from worker_config import app
from celery.result import AsyncResult

task_id = "a1b2c3d4-e5f6-7890-abcd-ef1234567890"
result_obj = AsyncResult(id=task_id, app=app)

if result_obj.ready():
    try:
        data = result_obj.get(timeout=10)
        print(data)
    except Exception as e:
        print(f'Retrieval Failed: {e}')

Scheduler Configuration for Periodic Work

For recurring operations, configure the Celery Beat scheduler alongside the worker process.

# beat_worker.py
import os
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "project.settings")

from celery import Celery
from celery.schedules import crontab
from datetime import timedelta

app = Celery(
    'scheduler_service', 
    include=['periodic_jobs'],
    broker='redis://127.0.0.1:6379/0'
)

app.conf.timezone = 'Asia/Shanghai'
app.conf.beat_schedule = {
    'daily_sync_task': {
        'task': 'periodic_jobs.refresh_content_cache',
        'schedule': timedelta(minutes=10),
        'args': (),
    },
}

Define the background job logic:

# periodic_jobs.py
from .beat_worker import app
from project.models import CourseItem
from django.core.cache import cache
from django.conf import settings

@app.task
def refresh_content_cache():
    items = CourseItem.objects.filter(is_active=True)[:settings.COURSE_LIMIT]
    
    cached_payload = [item.serialize() for item in items]
    
    # Persist to Redis
    cache.set('main_feed_cache', cached_payload, timeout=3600)
    return True

Start components separately:

  1. Worker: celery -A beat_worker worker -l info -P gevent
  2. Scheduler: celery -A beat_worker beat -l info

Vertical Database Partitioning

Large monolithic tables often benefit from splitting into domain-specific entities based on usage patterns. For instance, separating product categories (Standard, Premium, Lite) into distinct tables reduces locking contention and allows targeted indexing strategies.

Instead of a single wide table containing conditional flags for every category type, utilize an Abstract Base Modeel to share common attributes like created_at, name, or price. Each child table holds only unique properties relevant to its specific category.

This schema design improves query performance for specific segments while maintaining logical cohesion through shared parent schemas.

Tags: python Celery Database Design Async Processing Django

Posted on Tue, 06 Oct 2026 16:25:54 +0000 by hypertech