-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathprogress.py
More file actions
56 lines (38 loc) · 1.64 KB
/
Copy pathprogress.py
File metadata and controls
56 lines (38 loc) · 1.64 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
"""
Example of task progress reporting: a task updates its own progress via ProgressTracker while it runs, and the client
polls it through MongoResultBackend while the worker is still executing the task.
Requires a MongoDB instance, e.g. `make run_infra` from the repo root.
How to run:
1. Run worker: taskiq worker examples.progress:broker -w 1
2. Run client: uv run examples/progress.py
"""
import asyncio
import os
from taskiq import TaskiqDepends
from taskiq.depends.progress_tracker import ProgressTracker, TaskState
from taskiq_mongodb import MongoBroker, MongoResultBackend
MONGO_URI = os.environ.get("TASKIQ_MONGODB_URI", "mongodb://root:password@localhost:27017")
broker = MongoBroker(MONGO_URI, "taskiq_example").with_result_backend(
MongoResultBackend(MONGO_URI, "taskiq_example"),
)
@broker.task
async def process_batch(total: int, tracker: ProgressTracker[int] = TaskiqDepends()) -> str:
for done in range(1, total + 1):
await asyncio.sleep(0.5)
await tracker.set_progress(TaskState.STARTED, meta=done)
return f"processed {total} items"
async def main() -> None:
await broker.startup()
task = await process_batch.kiq(total=5)
last_meta = None
while not await task.is_ready():
progress = await task.get_progress()
if progress is not None and progress.meta != last_meta:
last_meta = progress.meta
print(f"progress: state={progress.state} meta={progress.meta}")
await asyncio.sleep(0.2)
result = await task.get_result()
print(f"result: {result.return_value}")
await broker.shutdown()
if __name__ == "__main__":
asyncio.run(main())