Nested Task (Submit A Parallel Task Within A Parallel Task)#
Nested tasks let tasks submit additional tasks without building a full graph up front.
"""Basic nested submit example (recursive Fibonacci)."""
from scaler import Client
from scaler.cluster.combo import SchedulerClusterCombo
def fibonacci(client: Client, n: int):
if n == 0:
return 0
if n == 1:
return 1
a = client.submit(fibonacci, client, n - 1)
b = client.submit(fibonacci, client, n - 2)
return a.result() + b.result()
def main():
cluster = SchedulerClusterCombo(n_workers=1)
with Client(address=cluster.get_address()) as client:
result = client.submit(fibonacci, client, 8).result()
print(result)
cluster.shutdown()
if __name__ == "__main__":
main()
What the example does:
Submits
fibonaccias a remote task.Inside each task, recursively submits child tasks with
submit().Waits on child futures and combines their results.
This pattern is useful to demonstrate nested execution, but recursion creates many small tasks and can be expensive.
Nested client addresses#
A nested Client given neither address takes both from the worker context.
address: the scheduler address its worker is connected to.object_storage_address: the address its worker reaches object storage on.
Pass address to reach another scheduler. Object storage is then the address that scheduler
advertises. Pass object_storage_address to set it directly.
This matters where the addresses a worker uses differ from the ones outside the cluster, for example behind NAT or a load balancer. See Object storage addresses.
from scaler import Client
def nested_task():
# Takes the scheduler and object storage addresses from the worker
with Client() as client:
result = client.submit(lambda x: x * 2, 5).result()
return result
client = Client(address="tcp://127.0.0.1:2345")
future = client.submit(nested_task)
print(future.result()) # 10