Celery Integration
The Celery integration provides a task class that can be used to automatically create and update tasks in Task Badger from Celery tasks. Although you can use the basic SDK functions to create and update tasks from Celery tasks, the Celery integration simplifies the usage significantly.
There are two ways you can use the Celery integration:
- Use the
CelerySystemIntegrationto automatically track all Celery tasks. - Use the
taskbadger.celery.Taskclass as the base class for Celery tasks you wish to track.
You can use both mechanisms at the same time since the base class is useful if you want to access to the Task Badger task object within the body of the Celery task.
Celery System Integration
If you want to track all tasks, you can use the CelerySystemIntegration class. By default, this will track every
task that is executed by the Celery workers (except the internal Celery tasks), including periodic / scheduled tasks.
import taskbadger
from taskbadger.systems.celery import CelerySystemIntegration
taskbadger.init(
token="YOUR_API_KEY",
systems=[CelerySystemIntegration()],
tags={"environment": "production"}
)
System Integration Options
The CelerySystemIntegration class takes a number of optional parameters:
auto_track_tasks: Set this toFalseto disable automatic tracking of tasks.includes: A list of task names or patterns to include. If this is set, only tasks that match one of the patterns will be tracked.excludes: A list of task names or patterns to exclude. If this is set, tasks that match one of the patterns will not be tracked.-
record_task_args: IfTrue, the arguments passed to the task will be recorded in the Task Badger task data.Since v1.4.0
-
heartbeat_interval: Seconds between automatic task updates while a task is running. See Keeping Long-Running Tasks Fresh.Since v2.4.0
-
stale_timeout: Thestale_timeoutto set on tracked tasks.Since v2.4.0
Exclusions take precedence over inclusions so if a task name matches both an include and an exclude, it will be excluded.
Celery Task Class
To track individual tasks, or if you want access to the Task object within the body of the Celery task, you can
use the taskbadger.celery.Task class as the base class for your Celery tasks. This can be used with or without
the CelerySystemIntegration.
This custom Celery task class tells Task Badger that the task should be tracked irrespective of any configuration
passed to CelerySystemIntegration ie. even if the task matches an exclusion rule it will still be tracked if it
is using taskbadger.celery.Task as its base.
The task class also provides convenient access to the Task Badger task object within the body of the Celery task.
Basic Usage
To use the integration simply set the base parameter of your Celery task to taskbadger.celery.Task:
This works the same with the @celery.shared_task decorator.
from celery import Celery
from taskbadger.celery import Task
app = Celery("tasks")
@app.task(base=Task)
def my_task():
pass
result = my_task.delay()
taskbadger_task_id = result.taskbadger_task_id
taskbadger_task = result.get_taskbadger_task()
Having made this change, a task will now be automatically created in Task Badger when the celery task is published. The Task Badger task will also be updated when the task completes.
Info
Note that Task Badger will only track the Celery task if it is run asynchronously. If the task is run
synchronously via .apply(), by calling the function directly, or if task_always_eager has been set,
the task will not be tracked.
This also means that the taskbadger_task_id attribute of the result as well as the return value
of result.get_taskbadger_task() will be None if the task is not being tracked by Task Badger.
Task Customization
You can pass additional parameters to the Task Badger Task class which will be used when creating the task.
This can be done by passing keyword arguments prefixed with taskbadger_ to the .apply_async() function or
to the task decorator.
taskbadger_ arguments to apply_async require the task base class
It is taskbadger.celery.Task that intercepts taskbadger_-prefixed arguments to apply_async, so
those only work on tasks declared with base=Task. On a plain Celery task they are silently ignored:
the task still publishes, and is still tracked if the
system integration tracks it, but the options have no effect.
taskbadger_ arguments on the task decorator are not affected. They are read off the task class
when the task is published, so they apply whether or not the task uses base=Task.
To set options per call on a task that is tracked by the system integration alone, pass them in the message headers instead — see Customization without the task base class.
# using the task decorator
@app.task(base=Task, taskbadger_monitor_id="xyz")
def my_task(arg1, arg2):
...
# using individual keyword arguments
my_task.apply_async(
args=[arg1, arg2],
taskbadger_name="my task",
taskbadger_value_max=1000,
taskbadger_data={"custom": "data"},
)
# using a dictionary
my_task.apply_async(args=[arg1, arg2], taskbadger_kwargs={
"name": "my task",
"value_max": 1000,
"data": {"custom": "data"}
})
Order of Precedence
Values passed via apply_async take precedence over values passed in the task decorator.
In both the decorator and apply_async, if individual keyword arguments are used as well as
the taskbadger_kwargs dictionary, the individual arguments will take precedence.
Recording task args
By default, the arguments passed to the task are not recorded in the Task Badger task data. To record the
arguments, set the taskbadger_record_task_args parameter to True in the task decorator or in the apply_async call.
This will override the value set in the CelerySystemIntegration if it is being used.
Since v1.4.0
Customization without the task base class
Task Badger options can also be passed in the Celery message headers, which works for any task,
including plain Celery tasks tracked by the CelerySystemIntegration:
my_task.apply_async(
args=[arg1, arg2],
headers={"taskbadger_kwargs": {
"name": "my task",
"value_max": 1000,
"data": {"custom": "data"},
}},
)
Unlike the taskbadger_-prefixed arguments, the header is read when the task is published rather than
by the task class, so no base=Task is needed. It takes the same options as taskbadger_kwargs, minus
the prefix, and takes precedence over values set on the task decorator.
Eager tasks read fewer options from the header
That applies to tasks which are actually published. When Celery runs a task eagerly
nothing is published, so the Task Badger task is created as the task starts instead, reading the
header directly. Only parent, heartbeat_interval and stale_timeout are honoured there —
name, value_max and data are ignored.
taskbadger_track is also required for an eager task that doesn't use base=Task, even if the
system integration would otherwise track it.
If the task would not otherwise be tracked — it isn't using the base class and doesn't match the system
integration's tracking rules — add taskbadger_track to the headers to track it anyway:
my_task.apply_async(
args=[arg1, arg2],
headers={
"taskbadger_track": True,
"taskbadger_kwargs": {"name": "my task"},
},
)
Note
record_task_args is a header of its own rather than an entry in taskbadger_kwargs:
headers={"taskbadger_record_task_args": True}.
The taskbadger_task_id attribute and get_taskbadger_task() method of the
result object are added by taskbadger.celery.Task, so they are not available on
the result when using headers alone.
Accessing the Task Object
The taskbadger.celery.Task class provides access to the Task Badger task object via the taskbadger_task property
of the Celery task. The Celery task instance can be accessed within a task function body by creating a
bound task.
@app.task(bind=True, base=taskbadger.celery.Task)
def my_task(self, items):
# Retrieve the Task Badger task
task = self.taskbadger_task
for i, item in enumerate(items):
do_something(item)
if i % 100 == 0:
# Track progress
task.update(value=i)
# Mark the task as complete
# This is normally handled automatically when the task completes but we call it here so that we
# can also update the `value` property or other task properties.
task.success(value=len(items))
Note
The taskbadger_task property will be None if the task is not being tracked by Task Badger.
This could indicate that the Task Badger API has not been configured, there was an error
creating the task, or the task is being run synchronously e.g. via .apply() or calling the task
using .map or .starmap, .chunk.
Keeping Long-Running Tasks Fresh
Since v2.4.0
A task with a stale_timeout is marked stale if it goes too long
without an update, so a long-running task that doesn't report progress will trip the timeout while it
is perfectly healthy. Setting heartbeat_interval (seconds) makes the worker update the task for you
while it runs, instead of having to do it from the task body.
The interval can be set on the system integration, on the task, or per call:
# for all tracked tasks
taskbadger.init(
token="YOUR_API_KEY",
systems=[CelerySystemIntegration(heartbeat_interval=60)],
)
# on the task
@app.task(base=Task, taskbadger_heartbeat_interval=60)
def my_task():
...
# per call
my_task.apply_async(taskbadger_heartbeat_interval=60)
Unless stale_timeout is given explicitly it is set to twice the interval, so each of the examples
above creates the task with a stale_timeout of 120 seconds. Pass both to control it:
As with the other options, values set on the task or on apply_async take precedence over the values
set on CelerySystemIntegration.
Note
All running tasks are updated from a single background thread per worker process, started the first time a task with a heartbeat runs. Updates stop when the task finishes.
Canvas primitives (map / starmap / chunks)
As of v1.6.3, Task Badger now tracks tasks created via Celery canvas primitives: map, starmap, and chunks. Previously these were executed as built-in celery.map / celery.starmap tasks and were filtered out; TaskBadger now creates task records for the inner tasks produced by these primitives.
Behavior summary:
- Canvas primitives produce Task Badger tasks using the inner task's name (so the tracked task has the same name as if the task had been executed individually).
- Each created task includes additional metadata:
canvas_type: one ofcelery.map,celery.starmap, orcelery.chunksitem_count: number of inner items produced by the canvas execution (an integer)celery_task_items: the arguments passed to the inner task
Example
Assume a Celery task:
Running a map:
Task Badger will create task records for each inner invocation with metadata similar to:
{
"name": "myapp.add",
"metadata": {
"canvas_type": "celery.map",
"item_count": 3,
"celery_task_items": [[1, 2], [2, 3], [3, 4]]
}
}
Subtasks
Since v2.5.0
A Celery task published from inside a tracked task is automatically nested under it via the
parent field, so you can see the work a task spawned.
Tasks nest a single level deep. A task published by a task that is itself a child becomes a sibling of that child rather than a grandchild.
What gets nested:
- Tasks published from the body of a tracked task via
.delay()or.apply_async(), including when Celery runs them eagerly. - Retries, which are nested under the first attempt.
What does not get nested:
- The next link of a
chain, andlinkcallbacks. These are successors of the task rather than work it chose to enqueue. - Tasks produced by canvas primitives (
map,starmap,chunks) when they run on a worker. Their Task Badger tasks are created in the worker rather than at publish time, so the enclosing task isn't visible there. They do nest when Celery runs eagerly.
Note
Known edge case: if a chain's next link is also called directly from the task body, that direct call is not nested.
To override the automatic nesting, pass taskbadger_parent for the call. An explicit None makes the
task a root task even though it was published from inside a tracked task:
# nest under a different task
my_task.apply_async(args=[arg1, arg2], taskbadger_parent=parent_task.id)
# opt out of nesting
my_task.apply_async(args=[arg1, arg2], taskbadger_parent=None)
taskbadger_parent takes a Task Badger task ID, not a Celery task or result ID.
As with the other taskbadger_ arguments to apply_async this requires the task to use base=Task;
without it, pass headers={"taskbadger_kwargs": {"parent": None}} instead. See
Customization without the task base class.
Since v2.5.1 eager tasks honour taskbadger_parent; before that they always nested under the
enclosing task.
Canvas primitives ignore taskbadger_parent
Canvas primitives (map, starmap, chunks) are
published under Celery's own celery.map / celery.starmap task rather than your own, so
taskbadger_parent is never intercepted and the taskbadger_kwargs header doesn't reach the
worker. On a worker their parent can't be set per call, and they are not nested at all.
Run eagerly they do nest under the enclosing task, and there the parent can be overridden with
headers={"taskbadger_kwargs": {"parent": ...}}.
External ID
The Celery task ID is automatically recorded on the Task Badger task's
external_id field, letting you correlate the Task Badger task with the
originating Celery task.
Opting out
To stop Task Badger tracking a single execution, set taskbadger_track to False, either as an
argument to apply_async or in the message headers:
# as an argument — needs base=Task
my_task.apply_async(args=[arg1, arg2], taskbadger_track=False)
# in the headers — works for any task
my_task.apply_async(args=[arg1, arg2], headers={"taskbadger_track": False})
# canvas primitives take the header form only
add.map([(1, 2), (2, 3)]).apply_async(headers={"taskbadger_track": False})
This takes precedence over the CelerySystemIntegration, so it opts out even when auto_track_tasks
is on, and over the base=Task class.
The value must be exactly False. Omitting it leaves tracking to the normal rules.
The argument form needs the task base class
As with the other taskbadger_ arguments to apply_async,
taskbadger_track=False is intercepted by taskbadger.celery.Task, so it only works on tasks
declared with base=Task. On a plain Celery task, or on a
canvas primitive, it is silently ignored and the task
stays tracked.
The header form works in all three cases, so prefer it unless you know the task uses base=Task.
To exclude a task from tracking on every call rather than per execution, use the excludes argument to
CelerySystemIntegration, which matches on task name.
Only works from v2.5.2
Since v2.5.2
Before that this header could only ever enable tracking, never suppress it. False was
indistinguishable from omitting the header, and apply_async on a base=Task task overwrote it
with True. On earlier versions use excludes instead.