acoustic-carpenter-78188
03/09/2023, 4:14 AM@task
def t1() -> Annotated[List[Any], 100]
Type
☐ Bug Fix
☑︎ Feature
☐ Plugin
Are all requirements met?
☐ Code completed
☐ Smoke tested
☐ Unit tests added
☐ Code documentation added
☐ Any pending items have an associated Issue
Complete description
Originally, in ListTransformer.to_literal, Items are converted one by one. It is time-consuming, especially for large lists.
flytekit/flytekit/core/type_engine.py
Line 977 in </flyteorg/flytekit/commit/8ccb5dde044f9dce3e216a8cea58dd464fb3ab6b|8ccb5dd>
This PR checks if python_type is list[flytePickle]. If so, batch the list of items and call FlytePickleTransformer.to_literal.
Measure the performance
Test file:
from typing import List, Dict, Any
from flytekit import task, Resources, workflow
@task(
limits=Resources(mem="4Gi",cpu="1"),
disable_deck=True,
)
def test() -> List[Any]:
return [{"a": {0: "foo"}}] * 10000
@workflow
def wf():
test()
if __name__ == "__main__":
wf()
# yichenglu/flytekit:mytest5 contains this PR
# pyflyte run --image yichenglu/flytekit:mytest5 --remote ./test.py wf
# <http://cr.flyte.org/flyteorg/flytekit:py3.10-1.4.0b1|cr.flyte.org/flyteorg/flytekit:py3.10-1.4.0b1> is the default image
# pyflyte run --image <http://cr.flyte.org/flyteorg/flytekit:py3.10-1.4.0b1|cr.flyte.org/flyteorg/flytekit:py3.10-1.4.0b1> --remote ./test.py wf
In an 4VCPU, 16G RAM EC2 instance, measuring the performance in flyte cluster using the above test file and the commands:
• Prior to these changes, , converting a list with type List[Any] with 10000 items needs 157 seconds.
• After this PR, convertinga list with type List[Dict[str, Any]] with 10000 items only needs 7 seconds.
image▾
image▾
from typing import List, Dict, Any, Annotated
from flytekit import task, Resources, workflow
@task(
limits=Resources(mem="4Gi",cpu="1"),
disable_deck=True,
container_image="{{.image.oldImage}}",
cache=True, cache_version="1.0"
)
def task0() -> List[List[Any]]:
return [ ["foo","foo","foo"] ] * 2
@task(
limits=Resources(mem="4Gi",cpu="1"),
disable_deck=True,
container_image="{{.image.prImage}}",
)
def task1(data: List[List[Any]]) -> List[List[Any]]:
print(data)
return data
@workflow
def wf():
data = task0()
task1(data=data)
# pyflyte run --image prImage="yichenglu/flytekit:mytest4" --image oldImage="<http://cr.flyte.org/flyteorg/flytekit:py3.10-1.4.1|cr.flyte.org/flyteorg/flytekit:py3.10-1.4.1>" --remote ./test.py wf
image▾
from typing import List, Dict, Any, Annotated
from flytekit import task, Resources, workflow
@task(
limits=Resources(mem="4Gi",cpu="1"),
disable_deck=True,
container_image="{{.image.prImage}}",
)
def task0() -> List[List[Any]]:
return [ ["foo","foo","foo"] ] * 2
@task(
limits=Resources(mem="4Gi",cpu="1"),
disable_deck=True,
container_image="{{.image.oldImage}}"
)
def task1(data: List[List[Any]]) -> List[List[Any]]:
print(data)
return data
@workflow
def wf():
data = task0()
task1(data=data)
# pyflyte run --image prImage="yichenglu/flytekit:mytest4" --image oldImage="<http://cr.flyte.org/flyteorg/flytekit:py3.10-1.4.1|cr.flyte.org/flyteorg/flytekit:py3.10-1.4.1>" --remote ./test.py wf
image▾
[[['foo', 'foo', 'foo']], [['foo', 'foo', 'foo']]] will be printed if see the pod's log
This occurs because in Task0, the outputs are batched, while in Task 1, which does not include this PR, the results are not expanded.
While slightly impacting backward compatibility, it is justifiable to ask users to update their images for two primary reasons:
• Before the implementation of this PR, users faced lengthy wait times for tasks involving pickling to be completed.
• Only tasks with the type list[any] will be affected. In most scenarios, users can still utilize the old images or cached outputs.
Tracking Issue
flyteorg/flyte#3207
Follow-up issue
NA
flyteorg/flytekit
✅ All checks have passed
30/30 successful checksacoustic-carpenter-78188
03/09/2023, 4:14 AMacoustic-carpenter-78188
04/04/2023, 9:19 PM