Skip to content

Tasks — QgsTaskManager

Tasks — qgis_sdk.tasks

Celery-like wrapper around QGIS background tasks that keeps UI responsive.

QGIS provides QgsApplication.taskManager() — global task manager handling progress, cancellation, dependencies, and UI feedback. This module wraps it with Celery-like API (@task, delay(), apply_async(), get(), ready(), AsyncResult, chain, group) and falls back to ThreadPoolExecutor for testing without QGIS.

from qgis_sdk.tasks import task
@task(bind=True)
def add(self, x, y):
self.set_progress(50)
return x + y
# Sync — like calling function directly
result = add(4, 4) # 8
# Async — like Celery delay
async_result = add.delay(4, 4)
print(async_result.get()) # 8
print(async_result.ready(), async_result.state) # True, SUCCESS

Installation

Terminal window
pip install qgis-sdk

Quick Start — Celery-like

from qgis_sdk.tasks import task, TaskManager, chain, group
from time import sleep
@task("My heavy work", bind=True, can_cancel=True)
def do_heavy_work(self, wait_time):
for i in range(100):
sleep(wait_time / 100.0)
self.set_progress(i)
if self.is_canceled():
return None
return {"result": 42}
def on_finished(exception, result):
if exception is None:
print(f"Task completed: {result}")
else:
print(f"Task failed: {exception}")
# Sync direct call — like Celery
result = do_heavy_work(2) # immediate
print(result)
# Async via delay — like Celery
async_result = do_heavy_work.delay(2)
print(async_result.get(timeout=10)) # waits
print(async_result.ready()) # True
print(async_result.successful()) # True
print(async_result.state) # SUCCESS, PENDING, FAILURE, REVOKED
print(async_result.id)
# apply_async with more control
async_result = do_heavy_work.apply_async(args=(2,), kwargs={}, countdown=1,
description="Custom", on_finished=on_finished)
# Signature / chain — like Celery
sig = do_heavy_work.s(2)
result = sig.delay()
print(result.get())
sig2 = do_heavy_work.si(2) # immutable
result = sig2.delay()
# Chain
from qgis_sdk.tasks import task
@task
def add(x, y=0):
return x + y
@task
def mul(x, y=1):
return x * y
c = chain(add.s(2, 2), add.s(3)) # add(2,2)=4, then add(4,3)=7
print(c()) # 7
print(c.delay().get()) # 7
# Group
results = group(add.s(1, 1), add.s(2, 2), add.s(3, 3))
print([r.get() for r in results]) # [2, 4, 6]
# Old API still works
from qgis_sdk.tasks import Task
t = Task.from_function("My task", do_heavy_work, wait_time=2, on_finished=on_finished, bind=True)
TaskManager.instance().add_task(t)

API

@task / @shared_task — Celery-like decorator

from qgis_sdk.tasks import task, shared_task, app
# With and without parentheses — like Celery
@task
def add(x, y):
return x + y
@task("My heavy work", bind=True, can_cancel=True, on_finished=on_finished)
def my_task(self, x, y):
self.set_progress(50)
if self.is_canceled():
return None
return x + y
@shared_task(bind=True)
def my_shared_task(self, x, y):
return x + y
# Via app — like Celery app
@app.task(bind=True)
def my_app_task(self, x, y):
return x + y
# Direct call — synchronous
result = add(4, 4) # 8, immediate
# Async — returns AsyncResult
async_result = add.delay(4, 4)
async_result = add.apply_async(args=(4, 4), kwargs={}, countdown=1, description="Custom")
# Signature
sig = add.s(4, 4)
sig = add.si(4, 4) # immutable
sig = add.signature(args=(4, 4), kwargs={}, immutable=False)
async_result = sig.delay()
result = sig() # sync via signature

Decorator returns TaskWrapper:

class TaskWrapper:
def __call__(self, *args, **kwargs) -> Any: ... # sync
def run(self, *args, **kwargs) -> Any: ... # alias
def delay(self, *args, **kwargs) -> AsyncResult: ...
def apply_async(self, args=None, kwargs=None, description=None, on_finished=None,
countdown=None, eta=None, bind=None, **options) -> AsyncResult: ...
def s(self, *args, **kwargs) -> Signature: ... # like Celery s()
def si(self, *args, **kwargs) -> Signature: ... # immutable
def signature(self, args=None, kwargs=None, immutable=False, **options) -> Signature: ...
@property
def name(self) -> str: ...
@property
def description(self) -> str: ...
  • bind=True: first arg is task instance (self) with set_progress, is_canceled, progress, etc. — like Celery bind=True.
  • If bind=False (default) but first param name is task/self, auto-bind is enabled for QGIS compatibility.
  • can_cancel, on_finished, flags passed to QGIS task.

AsyncResult — Celery-like

class AsyncResult:
@property
def id(self) -> str: ...
@property
def task_id(self) -> str: ... # alias
def get(self, timeout=None, propagate=True) -> Any: ... # waits, returns result, raises if failed and propagate=True
def wait(self, timeout=None, propagate=True) -> Any: ... # alias for get()
def ready(self) -> bool: ... # True if finished
def successful(self) -> bool: ... # True if finished and no exception
def failed(self) -> bool: ... # True if finished and exception
@property
def result(self) -> Any: ...
@property
def state(self) -> str: ... # PENDING, STARTED, SUCCESS, FAILURE, REVOKED
@property
def status(self) -> str: ... # alias for state
def revoke(self, terminate=False): ... # cancel task

Example:

result = my_task.delay(4, 4)
print(result.id)
print(result.get(timeout=10))
print(result.wait(timeout=10))
print(result.ready())
print(result.successful())
print(result.failed())
print(result.state)
print(result.result)
result.revoke()

Signature — Celery-like

class Signature:
def __init__(self, task_wrapper, args=(), kwargs=None, immutable=False, options=None): ...
def delay(self, *args, **kwargs) -> AsyncResult: ...
def apply_async(self, args=None, kwargs=None, **options) -> AsyncResult: ...
def __call__(self, *args, **kwargs) -> Any: ... # sync execution
def clone(self, args=None, kwargs=None, **opts) -> Signature: ...
def __or__(self, other) -> Chain: ... # chain via |
  • s(*args, **kwargs): mutable signature — chain passes previous result as first arg plus signature args.
  • si(*args, **kwargs): immutable signature — ignores previous result and new args in chain.

Chain and Group — Celery-like

from qgis_sdk.tasks import chain, group, task
@task
def add(x, y=0):
return x + y
@task
def mul(x, y=1):
return x * y
# Chain — previous result as first arg
c = chain(add.s(2, 2), add.s(3)) # 2+2=4, 4+3=7
print(c()) # 7
print(c.delay().get()) # 7
c2 = chain(add.s(2, 3), mul.s(2)) # 2+3=5, 5*2=10
print(c2()) # 10
# Via | operator — like Celery
c3 = add.s(2, 2) | add.s(3)
print(c3()) # 7
# Group — list of AsyncResults
results = group(add.s(1, 1), add.s(2, 2), add.s(3, 3))
print([r.get() for r in results]) # [2, 4, 6]

Task — QGIS wrapper (still available)

Wrapper around QgsTask.

class Task:
def __init__(self, description: str, flags=None, bind=False): ...
@classmethod
def from_function(cls, description: str, function: Callable, *args,
on_finished: Optional[Callable] = None, flags=None, bind=False, **kwargs) -> Task:
"""Create task from function.
function(task, *args, **kwargs) -> Any
on_finished(exception, result)
"""
def run(self) -> Any: ... # override in subclass
def finished(self, result: Any): ... # override
def set_progress(self, progress: float): ...
def setProgress(self, progress: float): ... # QGIS camelCase
def progress(self) -> float: ...
def is_canceled(self) -> bool: ...
def isCanceled(self) -> bool: ...
def is_finished(self) -> bool: ...
def isFinished(self) -> bool: ...
def cancel(self): ...
def can_cancel(self) -> bool: ...
def canCancel(self) -> bool: ...
@property
def id(self) -> str: ...
@property
def qgis_task(self): ... # underlying QgsTask if available
@property
def fallback_task(self): ... # fallback task

QGIS mapping:

  • Task.from_function → QgsTask.fromFunction
  • set_progress → setProgress
  • is_canceled → isCanceled
  • is_finished → status() == Complete or isFinished()
  • cancel → cancel()
  • can_cancel → canCancel()

Fallback:

  • Runs in ThreadPoolExecutor
  • set_progress stores value, calls callbacks
  • is_canceled checks flag
  • cancel sets flag and cancels future

TaskManager — QGIS global manager + Celery-like

Wrapper around QgsApplication.taskManager().

class TaskManager:
@classmethod
def instance(cls, max_workers: int = 4) -> TaskManager: ... # singleton
def add_task(self, task: Union[Task, Callable, TaskWrapper, Signature], *args, **kwargs) -> Union[Task, AsyncResult]:
"""Add task — accepts Task instance, callable, TaskWrapper, Signature.
If TaskWrapper/Signature, returns AsyncResult (celery-like).
"""
def tasks(self) -> List[Any]: ...
def count(self) -> int: ...
def cancel_all(self): ...
def shutdown(self, wait: bool = True): ...
@property
def qgis_manager(self): ... # QgsTaskManager if available
@property
def fallback_manager(self): ... # fallback manager

QGIS mapping:

  • instance() → QgsApplication.taskManager() singleton
  • add_task(task) → taskManager().addTask(task.qgis_task)
  • tasks() → taskManager().tasks()
  • cancel_all() → taskManager().cancelAll()

Fallback:

  • Uses ThreadPoolExecutor(max_workers=4)
  • add_task submits to executor, removes on done
  • count() returns running tasks

ProcessingAlgRunnerTask — with Celery-like delay

Wrapper around QgsProcessingAlgRunnerTask.

from qgis_sdk.tasks import ProcessingAlgRunnerTask, TaskManager
from qgis.core import QgsProcessingContext, QgsProcessingFeedback
context = QgsProcessingContext()
feedback = QgsProcessingFeedback()
params = {"INPUT": layer, "DISTANCE": 10, "OUTPUT": "memory:"}
task = ProcessingAlgRunnerTask("qgis:buffer", params, context, feedback,
on_executed=lambda successful, results: print(results))
# QGIS
task.qgis_task.executed.connect(lambda successful, results: print(results))
QgsApplication.taskManager().addTask(task.qgis_task)
# Celery-like
async_result = task.delay()
print(async_result.get())
async_result = task.apply_async()

Convenience Functions — now return AsyncResult

from qgis_sdk.tasks import run_task, add_task, cancel_all
async_result = run_task(my_function, wait_time=2, description="My task", on_finished=on_finished, bind=True)
print(async_result.get())
async_result = add_task(my_function, description="My task")
print(async_result.get())
async_result = add_task(my_task.s(4, 4)) # Signature
print(async_result.get())
cancel_all()

Inside QGIS Plugin — Celery-like

from qgis_sdk.tasks import task, TaskManager, chain, group
class MyPlugin:
@task("My task", bind=True, can_cancel=True)
def my_task(self, x, y):
self.set_progress(50)
if self.is_canceled():
return None
return x + y
def run_heavy(self):
# Sync
result = self.my_task(4, 4)
self.iface.messageBar().pushMessage(f"Result: {result}")
# Async via delay — like Celery
async_result = self.my_task.delay(4, 4)
def on_finished(exception, result):
if exception is None:
self.iface.messageBar().pushMessage(f"Task completed: {result}")
# For QGIS, on_finished is handled via Task.from_function
# For celery-like, you can also wait
# result = async_result.get(timeout=10)
self.iface.messageBar().pushMessage("Task started in background")
def run_chain(self):
# Chain tasks — like Celery
c = chain(self.my_task.s(2, 2), self.my_task.s(3))
result = c.delay()
# result.get() would be 7
def run_processing(self):
from qgis.core import QgsProcessingContext, QgsProcessingFeedback
from qgis_sdk.tasks import ProcessingAlgRunnerTask
from qgis.core import QgsApplication
context = QgsProcessingContext()
feedback = QgsProcessingFeedback()
params = {"INPUT": self.iface.activeLayer(), "DISTANCE": 10, "OUTPUT": "memory:"}
def on_executed(successful, results):
if successful:
self.iface.messageBar().pushMessage(f"Buffer completed")
task = ProcessingAlgRunnerTask("qgis:buffer", params, context, feedback, on_executed=on_executed)
QgsApplication.taskManager().addTask(task.qgis_task)
# Or celery-like
# async_result = task.delay()

Testing without QGIS — using own fixtures

conftest.py
pytest_plugins = ["qgis_sdk.testing"]
def test_task_manager(fake_task_manager):
# Uses fixture fake_task_manager — celery-like
def work(task):
task.set_progress(50)
return 42
async_result = fake_task_manager.add_task(work, description="Test", bind=True)
assert async_result.get() == 42
assert async_result.ready()
assert async_result.state == "SUCCESS"
assert fake_task_manager.count() == 0
def test_task_decorator_celery_like(fake_task_manager):
from qgis_sdk.tasks import task
@task(bind=True)
def add(self, x, y):
self.set_progress(50)
return x + y
# Sync
assert add(4, 4) == 8
# Async via delay — celery-like
async_result = add.delay(4, 4)
assert async_result.get() == 8
assert async_result.ready()
assert async_result.successful()
assert async_result.state == "SUCCESS"
assert async_result.id is not None
# apply_async
async_result2 = add.apply_async(args=(10,), kwargs={}, description="Custom")
assert async_result2.get() == 20
# Signature
sig = add.s(10, 20)
assert sig.delay().get() == 30
sig2 = add.si(5, 5)
assert sig2.delay().get() == 10
assert sig2.delay(100, 100).get() == 10 # immutable ignores new args
def test_chain_and_group():
from qgis_sdk.tasks import task, chain, group
@task
def add(x, y=0):
return x + y
c = chain(add.s(2, 2), add.s(3)) # 2+2=4, 4+3=7
assert c() == 7
results = group(add.s(1, 1), add.s(2, 2), add.s(3, 3))
assert [r.get() for r in results] == [2, 4, 6]
def test_fake_task_wrapper(fake_task_wrapper):
assert fake_task_wrapper(4, 4) == 8
assert fake_task_wrapper.delay(4, 4).get() == 8

FakeTask, FakeAsyncResult, FakeTaskManager, FakeTaskWrapper, FakeSignature:

from qgis_sdk.testing import FakeTask, FakeAsyncResult, FakeTaskManager, FakeTaskWrapper, FakeSignature
def work(task, value):
task.set_progress(100)
return value * 2
task = FakeTask("Test", work, value=10, bind=True)
assert task.description == "Test"
assert task.can_cancel() is True
task.set_progress(50)
assert task.progress() == 50
task._execute()
assert task.is_finished()
assert task.result() == 20
assert task.state == "SUCCESS"
# Celery-like delay
async_result = task.delay(value=20)
assert async_result.get() == 40
assert async_result.state == "SUCCESS"
assert async_result.ready()
assert async_result.successful()
# TaskWrapper
wrapper = FakeTaskWrapper(work, description="Test")
assert wrapper(10) == 20
assert wrapper.delay(10).get() == 20
assert wrapper.s(10).delay().get() == 20
# TaskManager
mgr = FakeTaskManager()
async_res = mgr.add_task(work, description="Test", bind=True)
assert async_res.get() == 42
assert mgr.count() == 0
assert len(mgr.added_tasks) == 1

Fallback Behavior

Without QGIS:

  • TaskManager uses ThreadPoolExecutor(max_workers=4)
  • Task.from_function creates _FallbackTask
  • TaskWrapper.delay() creates Task and adds to manager, returns AsyncResult
  • ProcessingAlgRunnerTask simulates with sleep 0.1s

With QGIS:

  • Uses QgsApplication.taskManager()
  • Task.from_function → QgsTask.fromFunction
  • TaskWrapper wraps QgsTask.fromFunction and returns AsyncResult
  • ProcessingAlgRunnerTask wraps QgsProcessingAlgRunnerTask

Comparison with QGIS Docs and Celery

From QGIS tasks cookbook:

# QGIS docs
from qgis.core import QgsTask, QgsApplication
def run(task, wait_time):
from time import sleep
for i in range(100):
sleep(wait_time / 100.0)
task.setProgress(i)
if task.isCanceled():
return None
return {"result": 42}
def finished(exception, result=None):
if exception is None:
print(f"Task completed: {result}")
task = QgsTask.fromFunction("My task", run, wait_time=2, on_finished=finished)
QgsApplication.taskManager().addTask(task)
# qgis_sdk — QGIS compatible
from qgis_sdk.tasks import Task, TaskManager
task = Task.from_function("My task", run, wait_time=2, on_finished=finished, bind=True)
TaskManager.instance().add_task(task)
# qgis_sdk — Celery-like (new)
from qgis_sdk.tasks import task
@task("My task", bind=True)
def run(self, wait_time):
for i in range(100):
self.set_progress(i)
if self.is_canceled():
return None
return {"result": 42}
result = run(2) # sync
async_result = run.delay(2) # async
print(async_result.get())

Celery docs:

# Celery
from celery import Celery
app = Celery('tasks', broker='pyamqp://guest@localhost//')
@app.task(bind=True)
def add(self, x, y):
return x + y
result = add.delay(4, 4)
print(result.get())
print(result.ready())
print(result.state)
# qgis_sdk — same API, but uses QGIS task manager or ThreadPoolExecutor
from qgis_sdk.tasks import task
@task(bind=True)
def add(self, x, y):
self.set_progress(50)
return x + y
result = add.delay(4, 4)
print(result.get())
print(result.ready())
print(result.state)

See Also