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 automatically transforms the target's interface to handle list-based inputs and outputs during workflow compilation and local execution.
Attributes
| Attribute | Type | Description |
|---|
| target | Union[LaunchPlan, ReferenceTask, FlyteLaunchPlan] | The target Flyte entity to map over |
| id | string | Unique identifier for the node, derived from the target entity's name. |
| metadata | _workflow_model.NodeMetadata = null | The metadata for the underlying node |
| concurrency | integer = null | If specified, this limits the number of mapped tasks than can run in parallel to the given batch size. |
| min_successes | integer = null | The minimum number of successful executions. If set, this takes precedence over min_success_ratio |
| min_success_ratio | float = 1.0 | The minimum ratio of successful executions. |
| bindings | List[_literal_models.Binding] = [] | A list of input bindings that define how data is passed to the mapped tasks. |
| python_interface | flyte_interface.Interface | The transformed Python interface representing the collection-based inputs and outputs of the array node. |
| interface | _interface_models.TypedInterface | The transformed system-level typed interface used for serialization and remote execution. |
| data_mode | _core_workflow.ArrayNode.DataMode | Determines how input data is partitioned, such as using a single input file or individual files per map instance. |
| execution_mode | _core_workflow.ArrayNode.ExecutionMode | Defines the state management strategy for the execution, such as full state or minimal state. |
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
| Name | Type | Description |
|---|
| target | Union[LaunchPlan, ReferenceTask, FlyteLaunchPlan] | The target Flyte entity to map over. |
| bindings | Optional[List[_literal_models.Binding]] = None | A list of input bindings for the node. |
| concurrency | Optional[int] = None | Limits the number of mapped tasks that can run in parallel. Set to 0 for unbounded concurrency. |
| min_successes | Optional[int] = None | The minimum number of successful executions required. Takes precedence over min_success_ratio. |
| min_success_ratio | Optional[float] = None | The minimum ratio of successful executions required (defaults to 1.0 if min_successes is not set). |
| metadata | Optional[_workflow_model.NodeMetadata] = None | The metadata for the underlying node. |
Methods
def construct_node_metadata(self) -> _workflow_model.NodeMetadata: ...
Constructs and returns the metadata for the node, defaulting to the target entity's name if no specific metadata is provided.
Returns
| Type | Description |
|---|
_workflow_model.NodeMetadata | The 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
| Type | Description |
|---|
str | The identifier string of the target entity. |
python_interface()
@property
def python_interface(self) -> flyte_interface.Interface: ...
Provides the Python-native interface definition for the array node, which typically involves list-transformed inputs and outputs.
Returns
| Type | Description |
|---|
flyte_interface.Interface | The Python interface object defining the expected input and output types. |
interface()
@property
def interface(self) -> _interface_models.TypedInterface: ...
Retrieves the typed interface for remote entities; raises an AttributeError if the interface is not available.
Returns
| Type | Description |
|---|
_interface_models.TypedInterface | The serialized interface model for the remote entity. |
bindings()
@property
def bindings(self) -> List[_literal_models.Binding]: ...
Returns the list of input bindings that map workflow data to the node's inputs.
Returns
| Type | Description |
|---|
List[_literal_models.Binding] | A list of binding models representing input connections. |
upstream_nodes()
@property
def upstream_nodes(self) -> List[Node]: ...
Returns a list of nodes that must execute before this node; currently returns an empty list for ArrayNodes.
Returns
| Type | Description |
|---|
List[Node] | An empty list as upstream dependencies are managed differently for array nodes. |
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
| Type | Description |
|---|
Any | The target Flyte entity object. |
data_mode()
@property
def data_mode(self) -> _core_workflow.ArrayNode.DataMode: ...
Indicates how data is handled for the array node, such as whether it uses single or individual input files.
Returns
| Type | Description |
|---|
_core_workflow.ArrayNode.DataMode | The data mode enum value determining input file structure. |
local_execute()
def local_execute(self, ctx: FlyteContext, **kwargs) -> Union[Tuple[Promise], Promise, VoidPromise]: ...
Executes the array node locally by iterating over input collections and invoking the target entity for each element. It validates success ratios and handles translation between Python types and Flyte literals.
Parameters
| Name | Type | Description |
|---|
| ctx | FlyteContext | The execution context providing environment configuration and state. |
Returns
| Type | Description |
|---|
Union[Tuple[Promise], Promise, VoidPromise] | A promise containing a collection of results, or a void promise if no outputs are expected. |
local_execution_mode()
def local_execution_mode(self): ...
Returns the execution mode for local processing.
Returns
| Type | Description |
|---|
ExecutionState.Mode | The mode indicating local task execution. |
min_success_ratio()
@property
def min_success_ratio(self) -> Optional[float]: ...
Returns the minimum ratio of successful sub-task executions required for the node to be considered successful.
Returns
| Type | Description |
|---|
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]: ...
Returns the absolute minimum number of successful sub-task executions required.
Returns
| Type | Description |
|---|
Optional[int] | The integer count of required successes. |
concurrency()
@property
def concurrency(self) -> Optional[int]: ...
Returns the maximum number of sub-tasks allowed to run in parallel.
Returns
| Type | Description |
|---|
Optional[int] | The batch size for parallel execution, or None for default behavior. |
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
| Type | Description |
|---|
_core_workflow.ArrayNode.ExecutionMode | The execution mode enum value. |
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
| Type | Description |
|---|
bool | Always returns True for this implementation. |
@property
def bound_inputs(self) -> Set[str]: ...
Returns the set of input names that are already bound to specific values.
Returns
| Type | Description |
|---|
Set[str] | An empty set of strings representing bound input identifiers. |