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
| Attribute | Type | Description |
|---|
| target | Union[LaunchPlan, ReferenceTask, FlyteLaunchPlan] | The target Flyte entity to map over |
| id | string | The unique identifier for the node, derived from the target entity's name. |
| metadata | NodeMetadata | The metadata for the underlying node |
| concurrency | integer | If specified, this limits the number of mapped tasks than can run in parallel to the given batch size. |
| min_successes | integer | 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[Binding] = [] | A list of input bindings that define how data is passed to the mapped tasks. |
| python_interface | Interface | The transformed interface representing the collection-based inputs and outputs for local execution. |
| interface | TypedInterface | The transformed typed interface used for serialization and remote execution. |
| data_mode | ArrayNode.DataMode | Determines whether inputs are handled as a single file or individual files based on the target type. |
| execution_mode | ArrayNode.ExecutionMode | Defines 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
| 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 | Metadata for the underlying node. |
Methods
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
| Type | Description |
|---|
_workflow_model.NodeMetadata | The 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
| Type | Description |
|---|
str | The 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
| Type | Description |
|---|
flyte_interface.Interface | The 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
| Type | Description |
|---|
_interface_models.TypedInterface | The 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
| Type | Description |
|---|
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
| Type | Description |
|---|
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
| Type | Description |
|---|
Any | The 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
| Type | Description |
|---|
_core_workflow.ArrayNode.DataMode | The 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
| Name | Type | Description |
|---|
| ctx | FlyteContext | The execution context used to manage state and literal translation. |
Returns
| Type | Description |
|---|
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
| Type | Description |
|---|
null | The 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
| Type | Description |
|---|
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
| Type | Description |
|---|
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
| Type | Description |
|---|
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
| Type | Description |
|---|
_core_workflow.ArrayNode.ExecutionMode | The 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
| Type | Description |
|---|
bool | Always returns True in this implementation. |
@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
| Type | Description |
|---|
Set[str] | An empty set of input names. |