Running parallel flows from an Agent tool¶
The Parallel multi-tenant scoring tutorial shows how an external service can create one temporary Application instance for each tenant and run their scoring flows in parallel. In some cases, you want an agent to trigger this workflow instead of calling the API directly. For example, an operations agent can start scoring for a list of tenants after validating a request, then summarize the outcome of each run.
Wrapping the workflow in an Agent tool gives the agent a constrained operation: it can provide tenant identifiers, trigger the scenario, and receive a structured result without handling Application instances itself. The tool also keeps the orchestration logic in one reusable place.
An Agent tool is not the right entry point for every scoring request. Use a direct API integration for high-volume, scheduled, or fully deterministic calls. An agent waits for the tool invocation to complete, so do not use this pattern for flows that exceed your agent timeout or for flows where the agent does not need to decide on the request
Prerequisites¶
The Parallel multi-tenant scoring tutorial was completed, including the
PROJECT_MTS_TEMPLATEApplication and itsSCORE_TENANTscenario.Dataiku >= 13.4.
An LLM connection
Creating the Agent tool¶
Create a Custom Python tool.
Go to the GenAI menu, select Agent Tools, click + New Agent Tool.
Select Custom Python, choose a meaningful name like Score tenants, and click Create.
Defining the parallel scoring tool¶
The tool accepts an array of tenant identifiers.
For each identifier, it creates a temporary Application instance, sets the BigQuery dataset variable,
and runs the scoring scenario.
A ThreadPoolExecutor starts several tenant runs at once, up to the configured MAX_WORKERS limit.
In the Design tab, replace the default code with Code 1.
import dataiku
from dataiku.llm.agent_tools import BaseAgentTool
import json
from concurrent.futures import ThreadPoolExecutor, as_completed
APP_ID = "PROJECT_MTS_TEMPLATE"
SCENARIO_ID = "SCORE_TENANT"
MAX_WORKERS = 4
class ScoreTenantsTool(BaseAgentTool):
"""Run scoring for one or more authorized tenants."""
def set_config(self, config, plugin_config):
pass
def get_descriptor(self, tool):
return {
"description": (
"Run the tenant scoring flow for an approved list of tenant IDs. "
"Use only after the user confirms the tenants to score."
),
"inputSchema": {
"type": "object",
"properties": {
"tenant_ids": {
"type": "array",
"description": "Tenant IDs to score.",
"items": {"type": "string"},
"minItems": 1,
}
},
"required": ["tenant_ids"],
},
}
def score_tenant(self, tenant_id):
try:
client = dataiku.api_client()
app_template = client.get_app(APP_ID)
with app_template.create_temporary_instance() as app_instance:
project = app_instance.get_as_project()
project.update_variables({"tenant_bq_dataset": tenant_id})
scenario = project.get_scenario(SCENARIO_ID)
scenario_run = scenario.run_and_wait()
return {
"tenant_id": tenant_id,
"run_id": scenario_run.id,
"outcome": scenario_run.outcome,
}
except Exception:
return {
"tenant_id": tenant_id,
"outcome": "ERROR",
}
def invoke(self, input, trace):
tenant_ids = input["input"]["tenant_ids"]
max_workers = min(MAX_WORKERS, len(tenant_ids))
results = []
with trace.subspan("Running tenant scoring flows") as subspan:
with ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = [
executor.submit(self.score_tenant, tenant_id)
for tenant_id in tenant_ids
]
for future in as_completed(futures):
results.append(future.result())
subspan.outputs["results"] = results
return {"output": json.dumps(results)}
The trace.subspan(...) block records the duration and outcome of the parallel run as a separate span.
You can inspect it with the Traces plugin.
The tool returns one result per tenant.
A failed scenario is reported in its result and does not prevent the other submitted runs from completing.
If your tool needs to read a managed dataset created by the temporary instance,
do so inside the with app_template.create_temporary_instance() block.
Closing the instance deletes its managed datasets.
In this tutorial, predictions remain in external BigQuery storage,
so the tool returns only each scenario’s outcome.
Attention
Do not expose unrestricted tenant identifiers to an agent. Validate them against an allowlist or enforce access controls in the data connection before using this tool in production.
Controlling concurrency¶
Set MAX_WORKERS to a value that matches the available capacity of your Dataiku instance.
The Agent tool can start flows in parallel, but it does not increase the number of available jobs,
activities, connections, or compute resources.
The tool invokes the scenario synchronously.
Keep MAX_WORKERS low enough that the complete tool invocation stays within the timeout configured for your agent environment.
For long-running scoring jobs, prefer a direct API service that returns a request identifier
and lets the caller retrieve the outcome later.
Tip
To make the workflow asynchronous, use one tool to start the scenario and another to check its status. The first tool must return the project key and scenario run ID, then the agent must call the status tool until it receives a final outcome. Use a non-temporary Application instance so that the project remains available between tool calls, and delete it after the run. This pattern requires explicit agent instructions and state management.
Creating a Visual Agent¶
In the GenAI menu, select Agents, click + New Agent,
select Visual Agent, and click Create.
Select an LLM connection, then add the Score tenants tool.
Give the agent clear instructions on when to use it. For example: “Use the Score tenants tool only after the user has confirmed the list of tenants to score.”
The descriptor in Code 1 tells the agent
that the tool expects a tenant_ids array.
When a user requests scoring for tenant_001 and tenant_002,
the agent can invoke the tool with the following input.
{
"input": {
"tenant_ids": ["tenant_001", "tenant_002"]
},
"context": {
}
}
Testing the agent¶
In the agent’s chat panel, enter a request such as “Score tenant_001 and tenant_002” and click Run test.
Confirm that the agent calls Score tenants and reports one outcome per tenant.
Calling the agent from the LLM Mesh API¶
To call the Visual Agent from an external application, get the agent ID from the URL of its edit page, or from this code snippet Use it as the LLM identifier, as shown in Code 3.
import dataikuapi
HOST = "https://DSS_HOST"
API_KEY = "API_KEY"
PROJECT_KEY = "AGENT_PROJECT"
AGENT_ID = "agent:AGENT_ID"
client = dataikuapi.DSSClient(HOST, API_KEY)
project = client.get_project(PROJECT_KEY)
llm = project.get_llm(AGENT_ID)
completion = llm.new_completion()
completion.with_message("Score tenant_001 and tenant_002.")
response = completion.execute()
if response.success:
print(response.text)
Wrapping up¶
Congratulations! You can now let a Visual Agent trigger isolated scoring flows for several tenants in parallel. The Agent tool controls how the flow is called, while temporary Application instances keep each tenant run separate.
Reference documentation¶
Classes¶
|
A handle to interact with an application on the DSS instance. |
|
Handle on an instance of an app. |
|
Entry point for the DSS API client |
|
A handle to interact with a project on the DSS instance. |
|
A handle to interact with a scenario on the DSS instance. |
A handle containing basic info about a past run of a scenario. |
|
Variant of |
Functions¶
Build an API client for the current DSS instance. |
|
Create a new temporary instance of this application. |
|
|
Get a handle to interact with a specific app. |
Get a handle on the project corresponding to this application instance. |
|
|
Get a handle to interact with a specific LLM |
|
Get a handle to interact with a specific project. |
|
Get a handle to interact with a specific scenario |
|
Request a run of the scenario and wait the end of the run to complete. |
|
Updates a set of variables for this project |
