Commit a92c38d
authored
Explicit Queue Listening (#537)
By default, a process running DBOS dequeues from all declared queues.
Now, you can instead use `DBOS.listen_queues` to explicitly tell a
process running DBOS to only dequeue workflows from a specific set of
queues.
```python
DBOS.listen_queues(
queues: List[Queue]
)
```
For example, you can use this to manage heterogeneous workers, so
workers of type A only dequeue and execute workflows from queue A while
workers of type B only dequeue and execute workflows from queue B. For
example:
```python
cpu_queue = Queue("queue_one")
gpu_queue = Queue("queue_two")
if __name__ == "__main__":
config = ...
worker_type = ... # "cpu' or 'gpu'
DBOS(config=config)
if worker_type = "gpu":
# GPU workers will only dequeue and execute workflows from the GPU queue
DBOS.listen_queues([gpu_queue])
elif worker_type == "cpu":
# CPU workers will only dequeue and execute workflows from the CPU queue
DBOS.listen_queues([cpu_queue])
DBOS.launch()
```
Also add logging of all queues being listened to.1 parent 977f769 commit a92c38d
3 files changed
+84
-2
lines changed| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
326 | 326 | | |
327 | 327 | | |
328 | 328 | | |
| 329 | + | |
329 | 330 | | |
330 | 331 | | |
331 | 332 | | |
| |||
1576 | 1577 | | |
1577 | 1578 | | |
1578 | 1579 | | |
| 1580 | + | |
| 1581 | + | |
| 1582 | + | |
| 1583 | + | |
| 1584 | + | |
| 1585 | + | |
| 1586 | + | |
| 1587 | + | |
| 1588 | + | |
| 1589 | + | |
| 1590 | + | |
| 1591 | + | |
| 1592 | + | |
| 1593 | + | |
| 1594 | + | |
1579 | 1595 | | |
1580 | 1596 | | |
1581 | 1597 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
175 | 175 | | |
176 | 176 | | |
177 | 177 | | |
| 178 | + | |
| 179 | + | |
| 180 | + | |
| 181 | + | |
| 182 | + | |
| 183 | + | |
| 184 | + | |
| 185 | + | |
| 186 | + | |
| 187 | + | |
178 | 188 | | |
179 | | - | |
180 | | - | |
| 189 | + | |
| 190 | + | |
| 191 | + | |
| 192 | + | |
| 193 | + | |
| 194 | + | |
| 195 | + | |
| 196 | + | |
181 | 197 | | |
182 | 198 | | |
183 | 199 | | |
| |||
212 | 228 | | |
213 | 229 | | |
214 | 230 | | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
| 235 | + | |
| 236 | + | |
| 237 | + | |
| 238 | + | |
| 239 | + | |
| 240 | + | |
| 241 | + | |
| 242 | + | |
| 243 | + | |
| 244 | + | |
| 245 | + | |
| 246 | + | |
| 247 | + | |
| 248 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1752 | 1752 | | |
1753 | 1753 | | |
1754 | 1754 | | |
| 1755 | + | |
| 1756 | + | |
| 1757 | + | |
| 1758 | + | |
| 1759 | + | |
| 1760 | + | |
| 1761 | + | |
| 1762 | + | |
| 1763 | + | |
| 1764 | + | |
| 1765 | + | |
| 1766 | + | |
| 1767 | + | |
| 1768 | + | |
| 1769 | + | |
| 1770 | + | |
| 1771 | + | |
| 1772 | + | |
| 1773 | + | |
| 1774 | + | |
| 1775 | + | |
| 1776 | + | |
| 1777 | + | |
| 1778 | + | |
| 1779 | + | |
| 1780 | + | |
| 1781 | + | |
| 1782 | + | |
| 1783 | + | |
| 1784 | + | |
| 1785 | + | |
| 1786 | + | |
0 commit comments