Skip to main content

ArrayNode

This class represents a node that executes a target Flyte entity, such as a launch plan or task, over a collection of inputs in parallel. It manages execution parameters including concurrency limits, minimum success counts, and success ratios to control how the array of sub-tasks is processed. The class automatically transforms the target's interface to handle list-based inputs and outputs, supporting both local execution and remote workflow compilation.

Attributes

AttributeTypeDescription
targetUnion[LaunchPlan, ReferenceTask, "FlyteLaunchPlan"]The target Flyte entity to map over
idstringUnique identifier for the node, derived from the target entity's name.
metadataOptional[_workflow_model.NodeMetadata]The metadata for the underlying node
namestringThe name of the target entity, used for node identification within the workflow.
python_interfaceflyte_interface.InterfaceThe transformed interface representing the list-based inputs and outputs for local execution.
interface_interface_models.TypedInterfaceThe transformed typed interface used for remote execution and serialization.
bindingsList[_literal_models.Binding] = []A list of input bindings that define how data is passed to the mapped task.
upstream_nodesList[Node] = []A list of nodes that must complete before this node can execute, currently returning an empty list.
flyte_entityAnyThe underlying Flyte entity (task or launch plan) being executed by this node.
data_mode_core_workflow.ArrayNode.DataModeDetermines how input data is partitioned, such as single input file or individual files.
min_success_ratioOptional[float] = 1.0The minimum ratio of successful executions.
min_successesOptional[int] = 0The minimum number of successful executions. If set, this takes precedence over min_success_ratio
concurrencyOptional[int]If specified, this limits the number of mapped tasks than can run in parallel to the given batch size.
execution_mode_core_workflow.ArrayNode.ExecutionModeDefines the state management strategy for the array node, such as full or minimal state.
is_original_sub_node_interfaceboolean = trueA boolean flag indicating if the node uses the original sub-node interface.
bound_inputsSet[str] = set()A set of input names that are bound to specific values and not mapped over.

Constructor

Signature

def ArrayNode(
self,
target: Union[LaunchPlan, ReferenceTask, "FlyteLaunchPlan"],
bindings: Optional[List[_literal_models.Binding]] = None,
concurrency: Optional[int] = None,
min_successes: Optional[int] = None,
min_success_ratio: Optional[float] = None,
metadata: Optional[_workflow_model.NodeMetadata] = None,
): ...

Parameters

NameTypeDescription
targetUnion[LaunchPlan, ReferenceTask, FlyteLaunchPlan]The target Flyte entity to map over.
bindingsOptional[List[_literal_models.Binding]] = NoneA list of input bindings for the node.
concurrencyOptional[int] = NoneLimits the number of mapped tasks running in parallel. 0 means unbounded; None inherits from the workflow.
min_successesOptional[int] = NoneThe minimum number of successful executions required.
min_success_ratioOptional[float] = NoneThe minimum ratio of successful executions required (defaults to 1.0 if min_successes is not set).
metadataOptional[_workflow_model.NodeMetadata] = NoneMetadata for the underlying node.

Methods


construct_node_metadata()

def construct_node_metadata(self) -> _workflow_model.NodeMetadata: ...

Constructs the metadata for the underlying node, defaulting to the target entity name if no specific metadata is provided.

Returns

TypeDescription
_workflow_model.NodeMetadataThe metadata object containing configuration for the workflow node.

name()

@property
def name(self) -> str: ...

Retrieves the name of the target Flyte entity associated with this node.

Returns

TypeDescription
strThe string identifier of the target entity.

python_interface()

@property
def python_interface(self) -> flyte_interface.Interface: ...

Provides the Python-native interface definition for the array node, typically transformed to represent list inputs/outputs.

Returns

TypeDescription
flyte_interface.InterfaceThe interface object defining the Python types for inputs and outputs.

interface()

@property
def interface(self) -> _interface_models.TypedInterface: ...

Retrieves the typed interface models for the node, used primarily during serialization.

Returns

TypeDescription
_interface_models.TypedInterfaceThe low-level typed interface model for the remote entity.

bindings()

@property
def bindings(self) -> List[_literal_models.Binding]: ...

Returns the list of literal bindings that map workflow inputs to the node's parameters.

Returns

TypeDescription
List[_literal_models.Binding]A list of bindings representing the data flow into the node.

upstream_nodes()

@property
def upstream_nodes(self) -> List[Node]: ...

Identifies the nodes that must complete execution before this node can start.

Returns

TypeDescription
List[Node]An empty list as upstream dependencies are managed at the workflow level.

flyte_entity()

@property
def flyte_entity(self) -> Any: ...

Returns the underlying Flyte entity (e.g., LaunchPlan or Task) that this node maps over.

Returns

TypeDescription
AnyThe target entity instance.

data_mode()

@property
def data_mode(self) -> _core_workflow.ArrayNode.DataMode: ...

Indicates how data is partitioned and passed to the sub-tasks (e.g., single input file vs individual files).

Returns

TypeDescription
_core_workflow.ArrayNode.DataModeThe data mode enum value determining input handling.

local_execute()

def local_execute(self, ctx: FlyteContext, **kwargs) -> Union[Tuple[Promise], Promise, VoidPromise]: ...

Simulates the array execution locally by iterating over input lists and invoking the target entity for each element.

Parameters

NameTypeDescription
ctxFlyteContextThe execution context providing access to configuration and state.
kwargsAnyThe keyword arguments representing the input lists to be mapped over.

Returns

TypeDescription
Union[Tuple[Promise], Promise, VoidPromise]A promise containing a collection of results from the successful sub-task executions.

local_execution_mode()

def local_execution_mode(self): ...

Defines the execution state mode for local runs.

Returns

TypeDescription
AnyThe local task execution mode constant.

min_success_ratio()

@property
def min_success_ratio(self) -> Optional[float]: ...

Fetches the minimum ratio of successful sub-task executions required for the ArrayNode to be considered successful.

Returns

TypeDescription
Optional[float]A float between 0 and 1, or None if min_successes is used instead.

min_successes()

@property
def min_successes(self) -> Optional[int]: ...

Fetches the absolute minimum number of successful sub-task executions required.

Returns

TypeDescription
Optional[int]The integer count of required successes.

concurrency()

@property
def concurrency(self) -> Optional[int]: ...

Fetches the limit on the number of sub-tasks that can run in parallel.

Returns

TypeDescription
Optional[int]The maximum batch size for parallel execution, or None for workflow default.

execution_mode()

@property
def execution_mode(self) -> _core_workflow.ArrayNode.ExecutionMode: ...

Indicates the execution strategy, such as FULL_STATE or MINIMAL_STATE, based on the target entity type.

Returns

TypeDescription
_core_workflow.ArrayNode.ExecutionModeThe execution mode enum value.

is_original_sub_node_interface()

@property
def is_original_sub_node_interface(self) -> bool: ...

Indicates if the node uses the original sub-node interface definition.

Returns

TypeDescription
boolAlways returns True for ArrayNode instances.

bound_inputs()

@property
def bound_inputs(self) -> Set[str]: ...

Returns the set of input names that are bound to specific values rather than mapped over.

Returns

TypeDescription
Set[str]An empty set as bound inputs are currently not supported.