Skip to main content

A Celery Beat scheduler that stores the schedule in MongoDB.

Project description

Celery-MongoBeat

A modern, drop-in replacement for celerybeat-mongo. This project provides a Celery Beat scheduler that stores and retrieves task schedules from a MongoDB collection, allowing for dynamic management of periodic tasks without restarting the Celery Beat service.

Why celery-mongobeat?

The original celerybeat-mongo library is no longer actively maintained and contains several critical bugs. This project was created to provide a stable, reliable, and modern alternative for the community, ensuring continued support for dynamic, database-backed Celery schedules.

Features

  • Stable and Reliable: Fixes critical bugs from celerybeat-mongo, such as the issue where disabling one task would prevent all tasks from running.
  • Dynamic Task Management: Add, modify, and remove periodic tasks on the fly without restarting the beat service.
  • MongoDB Backend: Leverages MongoDB for a robust and scalable schedule store.
  • Fine-Grained Control:
    • Run Count Limiting: Use max_run_count to run a task a specific number of times and then automatically disable it.
  • Flexible Configuration: Full support for advanced pymongo.MongoClient options (like SSL) via mongodb_scheduler_client_kwargs.
  • Backwards Compatible: Supports legacy configuration variables from celerybeat-mongo for a smoother transition.
  • Modern Tooling: Built with a modern Python packaging structure (pyproject.toml).
  • All Schedule Types: Natively supports interval, crontab, and solar schedules.

Installation

Install the package from PyPI:

pip install celery-mongobeat

Configuration

To use this scheduler, set the beat_scheduler option in your Celery configuration.

Recommended Configuration

# celeryconfig.py

mongodb_scheduler_url = "mongodb://localhost:27017/"
mongodb_scheduler_db = "celery"
mongodb_scheduler_collection = "schedules"

beat_scheduler = "celery_mongobeat.beat:MongoScheduler"

Migrating from celerybeat-mongo

celery-mongobeat is designed as a near drop-in replacement, but there is one important configuration change you must make when migrating:

  • Update the Scheduler Path: The import path for the scheduler has been updated to align with modern package structures and Celery best practices.

You must change your beat_scheduler setting from: 'celerybeat_mongo.schedulers.MongoScheduler' (the old path) to: 'celery_mongobeat.beat:MongoScheduler' (the new path)


### Legacy (Backwards-Compatible) Configuration

If you are migrating from `celerybeat-mongo`, this library provides backward compatibility for the uppercase configuration variables. Modern, lowercase settings (e.g., `mongodb_scheduler_url`) will always take precedence.

```python
# celeryconfig.py

# Legacy uppercase individual settings (from celerybeat-mongo)
CELERY_MONGODB_SCHEDULER_URL = "mongodb://localhost:27017/"
CELERY_MONGODB_SCHEDULER_DB = "celery"
CELERY_MONGODB_SCHEDULER_COLLECTION = "schedules"

beat_scheduler = "celery_mongobeat.beat:MongoScheduler"

Usage

Once configured, start Celery Beat as you normally would:

celery -A your_app beat -l info

You can now manage your schedules by adding, updating, or removing documents in the configured MongoDB collection.

Programmatic Usage Example

For users who prefer a programmatic API over manually inserting documents into MongoDB, celery-mongobeat provides a convenient ScheduleManager helper class.

This allows you to easily create, update, and disable tasks from within your application code.

Example Usage

# In your application code (e.g., a management script or view)
from celery import current_app
from celery_mongobeat.helpers import ScheduleManager

# The recommended way to get a manager instance.
# It automatically reads the database configuration from your Celery settings.
app = current_app._get_current_object()
manager = ScheduleManager.from_celery_app(app)

# Example: Create a task to run every 30 seconds
manager.create_interval_task(
    name='my-periodic-task',
    task='your_project.tasks.some_task',
    every=30,
    period='seconds',
    args=[1, 2, 3]
)
print("Upserted interval task: 'my-periodic-task'")

# Example: Create a task with a description
manager.create_crontab_task(
    name='daily-report',
    task='your_project.tasks.generate_report',
    minute='0',
    hour='4',  # Run at 4:00 AM daily
    description='This is a custom description for the daily report task.'
)

# Example: Create a task that runs 5 times and then stops
manager.create_interval_task(
    name='run-five-times-task',
    task='your_project.tasks.some_task',
    every=60,
    period='seconds',
    max_run_count=5
)
print("Upserted limited-run task: 'run-five-times-task'")

# Example: Disable a task
manager.disable_task('my-periodic-task')
print("Disabled task: 'my-periodic-task'")

# Example: Get all enabled interval tasks
enabled_interval_tasks = manager.get_tasks(enabled=True, schedule_type='interval')
print(f"Found {len(enabled_interval_tasks)} enabled interval tasks.")
for task in enabled_interval_tasks:
    print(f" - {task['name']}")

### Creating Tasks from a Dictionary

The `create_*_task` methods are designed to be flexible. You can use Python's keyword argument unpacking (`**`) to create tasks from a dictionary. This is especially useful when processing data from an API or another data source.

Any extra keys in the dictionary that do not match a method parameter will be safely ignored.

```python
# Example data that might come from a web form or API
task_data = {
    'name': 'api-created-task',
    'task': 'your_project.tasks.process_data',
    'every': 15,
    'period': 'minutes',
    'args': [12345],
    'metadata': 'Created by API endpoint /tasks',  # This key will be ignored
    'request_id': 'xyz-789'  # This key will also be ignored
}

task_id = manager.create_interval_task(**task_data)
print(f"Successfully created task from dictionary with ID: {task_id}")

### Advanced Usage: Subclassing and Direct Database Access

The `ScheduleManager` is designed to be a flexible base. For more complex applications, it is highly recommended to subclass it to create a domain-specific API for your tasks. This encapsulates your application's scheduling logic, making your code cleaner and more maintainable.
 
**1. Subclassing `ScheduleManager`**

```python
# In your_app/scheduling.py

from celery_mongobeat.helpers import ScheduleManager

class AppScheduleManager(ScheduleManager):
    """A custom manager for our application's specific tasks."""

    def create_user_report_task(self, user_id: int):
        """Creates a recurring daily report for a specific user."""
        task_name = f"user-report-{user_id}"
        super().create_crontab_task(
            name=task_name,
            task='your_app.tasks.generate_report',
            kwargs={'user_id': user_id},
            minute='0',  # At the start of the hour
            hour='3'     # At 3 AM
        )
        print(f"Scheduled daily report for user {user_id}.")

# In your application code, you can now use this custom manager:
# from celery import current_app
# from your_app.scheduling import AppScheduleManager

# app = current_app._get_current_object()
# app_manager = AppScheduleManager.from_celery_app(app)
# app_manager.create_user_report_task(user_id=123)

2. Direct Database Access

For advanced queries, such as MongoDB aggregation pipelines, you can and should use the pymongo collection object that you used to initialize the manager. This gives you the full power of pymongo for any use case not directly covered by the helper.

# For example, to find the most common task paths using an aggregation:
pipeline = [
    {"$group": {"_id": "$task", "count": {"$sum": 1}}},
    {"$sort": {"count": -1}}
]
most_common_tasks = list(schedules_collection.aggregate(pipeline))
print("Most common tasks:", most_common_tasks)

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

celery_mongobeat-0.1.13.tar.gz (26.5 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

celery_mongobeat-0.1.13-py3-none-any.whl (20.1 kB view details)

Uploaded Python 3

File details

Details for the file celery_mongobeat-0.1.13.tar.gz.

File metadata

  • Download URL: celery_mongobeat-0.1.13.tar.gz
  • Upload date:
  • Size: 26.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.7

File hashes

Hashes for celery_mongobeat-0.1.13.tar.gz
Algorithm Hash digest
SHA256 3f16e40137ca94a1e606c7917a687456ead1f78d275033461e79ba244af3929c
MD5 a188520764c10f26cf7641a03ec775f6
BLAKE2b-256 27685a6fdae5884c78a0651c33e2af0a1ec3b08673cf9ab217b8acb24e68d1dc

See more details on using hashes here.

Provenance

The following attestation bundles were made for celery_mongobeat-0.1.13.tar.gz:

Publisher: publish.yml on sockmysox/celery-mongobeat

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file celery_mongobeat-0.1.13-py3-none-any.whl.

File metadata

File hashes

Hashes for celery_mongobeat-0.1.13-py3-none-any.whl
Algorithm Hash digest
SHA256 760fdf2fa8a7f080f9dceec10ece4e9eaf434a919b80f668c6c854d545a248c6
MD5 62bab6f6b9dcdf0a8f844c81109b05d8
BLAKE2b-256 4517e284d8c72c69c43264159481835e6add05739740eeb31077bce74510871d

See more details on using hashes here.

Provenance

The following attestation bundles were made for celery_mongobeat-0.1.13-py3-none-any.whl:

Publisher: publish.yml on sockmysox/celery-mongobeat

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page