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 provides configurable controls for execution concurrency and success thresholds, allowing for partial failures based on a minimum success count or ratio. The class manages the transformation of interfaces to handle list-based inputs and outputs while supporting both local execution and remote serialization.

Attributes

AttributeTypeDescription
targetUnion[LaunchPlan, ReferenceTask, FlyteLaunchPlan]The target Flyte entity to map over
idstringThe unique identifier for the node, derived from the target entity's name.
metadataNodeMetadataThe metadata for the underlying node
concurrencyintegerIf specified, this limits the number of mapped tasks than can run in parallel to the given batch size.
min_successesintegerThe minimum number of successful executions. If set, this takes precedence over min_success_ratio
min_success_ratiofloat = 1.0The minimum ratio of successful executions.
bindingsList[Binding] = []A list of input bindings that define how data is passed to the mapped tasks.
python_interfaceInterfaceThe transformed interface representing the collection-based inputs and outputs for local execution.
interfaceTypedInterfaceThe transformed typed interface used for serialization and remote execution.
data_modeArrayNode.DataModeDetermines whether inputs are handled as a single file or individual files based on the target type.
execution_modeArrayNode.ExecutionModeDefines the state management strategy (FULL_STATE or MINIMAL_STATE) used during execution.

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 that can run in parallel. Set to 0 for unbounded concurrency.
min_successesOptional[int] = NoneThe minimum number of successful executions required. Takes precedence over min_success_ratio.
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 node, defaulting to the target entity's name if no specific metadata is provided.

Returns

TypeDescription
_workflow_model.NodeMetadataThe metadata object containing configuration for the underlying Flyte node.

name()

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

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

Returns

TypeDescription
strThe name of the target entity.

python_interface()

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

Returns the Python-native interface for the ArrayNode, which is transformed to handle list-based inputs and outputs.

Returns

TypeDescription
flyte_interface.InterfaceThe transformed interface suitable for mapping over collections.

interface()

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

Retrieves the serialized typed interface for the node, typically used during remote execution or serialization.

Returns

TypeDescription
_interface_models.TypedInterfaceThe remote-compatible typed interface.

bindings()

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

Returns the list of input bindings that map workflow variables to the node's inputs.

Returns

TypeDescription
List[_literal_models.Binding]A list of literal bindings for the node.

upstream_nodes()

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

Returns the list of nodes that this node depends on; currently returns an empty list for ArrayNodes.

Returns

TypeDescription
List[Node]An empty list of upstream node dependencies.

flyte_entity()

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

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

Returns

TypeDescription
AnyThe target Flyte entity.

data_mode()

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

Indicates how data is transferred to the sub-nodes, such as using single or individual input files.

Returns

TypeDescription
_core_workflow.ArrayNode.DataModeThe data transfer mode for the array execution.

local_execute()

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

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

Parameters

NameTypeDescription
ctxFlyteContextThe execution context used to manage state and literal translation.

Returns

TypeDescription
Union[Tuple[Promise], Promise, VoidPromise]A promise containing a collection of results from the mapped executions.

local_execution_mode()

def local_execution_mode(self): ...

Returns the execution mode for local runs, identifying this as a local task execution.

Returns

TypeDescription
nullThe local task execution mode constant.

min_success_ratio()

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

Returns the minimum ratio of successful sub-node executions required for the ArrayNode to be considered successful.

Returns

TypeDescription
Optional[float]A float between 0 and 1 representing the success threshold ratio.

min_successes()

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

Returns the absolute minimum number of successful sub-node executions required.

Returns

TypeDescription
Optional[int]The integer count of required successful executions.

concurrency()

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

Returns the maximum number of sub-nodes allowed to run in parallel.

Returns

TypeDescription
Optional[int]The concurrency limit, or None if inheriting from the workflow.

execution_mode()

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

Indicates the execution strategy for the array node, such as FULL_STATE or MINIMAL_STATE.

Returns

TypeDescription
_core_workflow.ArrayNode.ExecutionModeThe execution mode constant.

is_original_sub_node_interface()

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

Returns a boolean indicating if the node uses the original sub-node interface.

Returns

TypeDescription
boolAlways returns True in this implementation.

bound_inputs()

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

Returns the set of input names that are bound to specific values; currently returns an empty set.

Returns

TypeDescription
Set[str]An empty set of input names.