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 directlyresult = add(4, 4) # 8
# Async — like Celery delayasync_result = add.delay(4, 4)print(async_result.get()) # 8print(async_result.ready(), async_result.state) # True, SUCCESSInstallation
pip install qgis-sdkQuick Start — Celery-like
from qgis_sdk.tasks import task, TaskManager, chain, groupfrom 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 Celeryresult = do_heavy_work(2) # immediateprint(result)
# Async via delay — like Celeryasync_result = do_heavy_work.delay(2)print(async_result.get(timeout=10)) # waitsprint(async_result.ready()) # Trueprint(async_result.successful()) # Trueprint(async_result.state) # SUCCESS, PENDING, FAILURE, REVOKEDprint(async_result.id)
# apply_async with more controlasync_result = do_heavy_work.apply_async(args=(2,), kwargs={}, countdown=1, description="Custom", on_finished=on_finished)
# Signature / chain — like Celerysig = do_heavy_work.s(2)result = sig.delay()print(result.get())
sig2 = do_heavy_work.si(2) # immutableresult = sig2.delay()
# Chainfrom qgis_sdk.tasks import task
@taskdef add(x, y=0): return x + y
@taskdef mul(x, y=1): return x * y
c = chain(add.s(2, 2), add.s(3)) # add(2,2)=4, then add(4,3)=7print(c()) # 7print(c.delay().get()) # 7
# Groupresults = 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 worksfrom qgis_sdk.tasks import Taskt = 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@taskdef 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 — synchronousresult = add(4, 4) # 8, immediate
# Async — returns AsyncResultasync_result = add.delay(4, 4)async_result = add.apply_async(args=(4, 4), kwargs={}, countdown=1, description="Custom")
# Signaturesig = add.s(4, 4)sig = add.si(4, 4) # immutablesig = add.signature(args=(4, 4), kwargs={}, immutable=False)async_result = sig.delay()result = sig() # sync via signatureDecorator 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) withset_progress,is_canceled,progress, etc. — like Celerybind=True.- If
bind=False(default) but first param name istask/self, auto-bind is enabled for QGIS compatibility. can_cancel,on_finished,flagspassed 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 taskExample:
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
@taskdef add(x, y=0): return x + y
@taskdef mul(x, y=1): return x * y
# Chain — previous result as first argc = chain(add.s(2, 2), add.s(3)) # 2+2=4, 4+3=7print(c()) # 7print(c.delay().get()) # 7
c2 = chain(add.s(2, 3), mul.s(2)) # 2+3=5, 5*2=10print(c2()) # 10
# Via | operator — like Celeryc3 = add.s(2, 2) | add.s(3)print(c3()) # 7
# Group — list of AsyncResultsresults = 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 taskQGIS mapping:
Task.from_function→QgsTask.fromFunctionset_progress→setProgressis_canceled→isCanceledis_finished→status() == CompleteorisFinished()cancel→cancel()can_cancel→canCancel()
Fallback:
- Runs in
ThreadPoolExecutor set_progressstores value, calls callbacksis_canceledchecks flagcancelsets 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 managerQGIS mapping:
instance()→QgsApplication.taskManager()singletonadd_task(task)→taskManager().addTask(task.qgis_task)tasks()→taskManager().tasks()cancel_all()→taskManager().cancelAll()
Fallback:
- Uses
ThreadPoolExecutor(max_workers=4) add_tasksubmits to executor, removes on donecount()returns running tasks
ProcessingAlgRunnerTask — with Celery-like delay
Wrapper around QgsProcessingAlgRunnerTask.
from qgis_sdk.tasks import ProcessingAlgRunnerTask, TaskManagerfrom 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))
# QGIStask.qgis_task.executed.connect(lambda successful, results: print(results))QgsApplication.taskManager().addTask(task.qgis_task)
# Celery-likeasync_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)) # Signatureprint(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
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() == 8FakeTask, 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() == 20assert task.state == "SUCCESS"
# Celery-like delayasync_result = task.delay(value=20)assert async_result.get() == 40assert async_result.state == "SUCCESS"assert async_result.ready()assert async_result.successful()
# TaskWrapperwrapper = FakeTaskWrapper(work, description="Test")assert wrapper(10) == 20assert wrapper.delay(10).get() == 20assert wrapper.s(10).delay().get() == 20
# TaskManagermgr = FakeTaskManager()async_res = mgr.add_task(work, description="Test", bind=True)assert async_res.get() == 42assert mgr.count() == 0assert len(mgr.added_tasks) == 1Fallback Behavior
Without QGIS:
TaskManagerusesThreadPoolExecutor(max_workers=4)Task.from_functioncreates_FallbackTaskTaskWrapper.delay()createsTaskand adds to manager, returnsAsyncResultProcessingAlgRunnerTasksimulates with sleep 0.1s
With QGIS:
- Uses
QgsApplication.taskManager() Task.from_function→QgsTask.fromFunctionTaskWrapperwrapsQgsTask.fromFunctionand returnsAsyncResultProcessingAlgRunnerTaskwrapsQgsProcessingAlgRunnerTask
Comparison with QGIS Docs and Celery
From QGIS tasks cookbook:
# QGIS docsfrom 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 compatiblefrom 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) # syncasync_result = run.delay(2) # asyncprint(async_result.get())Celery docs:
# Celeryfrom celery import Celeryapp = 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 ThreadPoolExecutorfrom 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)