acoustic-carpenter-78188
03/09/2023, 11:35 PMpyflyte-map-execute ....
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
image▾
import typing
from typing import List
from flytekit import map_task, task, workflow, ContainerTask, kwtypes
calculate_ellipse_area_shell = ContainerTask(
name="ellipse-area-metadata-python",
input_data_dir="/var/inputs",
output_data_dir="/var/outputs",
inputs=kwtypes(a=int),
outputs=kwtypes(area=float),
image="pingsutw/raw-container:v9",
command=[
"python",
"test.py",
"{{.inputs.a}}",
"/var/outputs",
],
)
@task
def coalesce(b: List[str]) -> str:
coalesced = "".join(b)
return coalesced
@task
def g_l(n: int) -> List[int]:
res = []
for i in range(n):
res.append(i)
return res
@workflow
def wf(n: int = 2):
l = g_l(n=n)
map_task(calculate_ellipse_area_shell)(a=l)
if __name__ == "__main__":
result = wf()
• dockerfile
FROM python:3.10-slim-buster
WORKDIR /root
COPY *.py /root/
• test.py
import math
import sys
import os
def write_output(output_dir, output_file, v):
with open(f"{output_dir}/{output_file}", "w") as f:
f.write(str(v))
def calculate_area(a, b):
return math.pi * a * b
def main(a, output_dir):
# parse list
li = a.strip('][').split(',')
res = [eval(i) for i in li]
# get the job index
index = int(os.environ.get(os.environ.get("BATCH_JOB_ARRAY_INDEX_VAR_NAME")))
area = calculate_area(res[index], 3)
write_output(output_dir, "area", f"{area}")
write_output(output_dir, "_SUCCESS", "")
if __name__ == "__main__":
a = sys.argv[1]
output_dir = sys.argv[2]
main(a, output_dir)
Tracking Issue
https://flyte-org.slack.com/archives/CP2HDHKE1/p1678230956906899
Follow-up issue
flyteorg/flytecopilot#54
flyteorg/flyteplugins#329
flyteorg/flytekit
✅ All checks have passed
30/30 successful checks