-
Notifications
You must be signed in to change notification settings - Fork 123
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Implementation of taskvine allpairs/map/reduce (#4011)
* Implementation of taskvine allpairs/map/reduce * lint * lint v2 * cleanup code * cleanup reduce * add test * remove debug print * cleanup map * format * allpairs in terms of map * format * do not create lib in map * error on lib name --------- Co-authored-by: Kevin Xue <[email protected]> Co-authored-by: Benjamin Tovar <[email protected]>
- Loading branch information
1 parent
8113f40
commit c1b9c8c
Showing
3 changed files
with
279 additions
and
22 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,84 @@ | ||
#!/bin/sh | ||
|
||
set -e | ||
|
||
. ../../dttools/test/test_runner_common.sh | ||
|
||
import_config_val CCTOOLS_PYTHON_TEST_EXEC | ||
import_config_val CCTOOLS_PYTHON_TEST_DIR | ||
|
||
export PYTHONPATH=$(pwd)/../../test_support/python_modules/${CCTOOLS_PYTHON_TEST_DIR}:$PYTHONPATH | ||
|
||
STATUS_FILE=vine.status | ||
PORT_FILE=vine.port | ||
|
||
check_needed() | ||
{ | ||
[ -n "${CCTOOLS_PYTHON_TEST_EXEC}" ] || return 1 | ||
|
||
# Poncho currently requires ast.unparse to serialize the function, | ||
# which only became available in Python 3.9. Some older platforms | ||
# (e.g. almalinux8) will not have this natively. | ||
"${CCTOOLS_PYTHON_TEST_EXEC}" -c "from ast import unparse" || return 1 | ||
|
||
# In some limited build circumstances (e.g. macos build on github), | ||
# poncho doesn't work due to lack of conda-pack or cloudpickle | ||
"${CCTOOLS_PYTHON_TEST_EXEC}" -c "import conda_pack" || return 1 | ||
"${CCTOOLS_PYTHON_TEST_EXEC}" -c "import cloudpickle" || return 1 | ||
|
||
return 0 | ||
} | ||
|
||
prepare() | ||
{ | ||
rm -f $STATUS_FILE | ||
rm -f $PORT_FILE | ||
return 0 | ||
} | ||
|
||
run() | ||
{ | ||
( ${CCTOOLS_PYTHON_TEST_EXEC} vine_python_future_hof.py $PORT_FILE; echo $? > $STATUS_FILE ) & | ||
|
||
# wait at most 15 seconds for vine to find a port. | ||
wait_for_file_creation $PORT_FILE 15 | ||
|
||
run_taskvine_worker $PORT_FILE worker.log --cores 2 --memory 2000 --disk 2000 | ||
|
||
# wait for vine to exit. | ||
wait_for_file_creation $STATUS_FILE 15 | ||
|
||
# retrieve exit status | ||
status=$(cat $STATUS_FILE) | ||
if [ $status -ne 0 ] | ||
then | ||
# display log files in case of failure. | ||
logfile=$(latest_vine_debug_log) | ||
if [ -f ${logfile} ] | ||
then | ||
echo "master log:" | ||
cat ${logfile} | ||
fi | ||
|
||
if [ -f worker.log ] | ||
then | ||
echo "worker log:" | ||
cat worker.log | ||
fi | ||
|
||
exit 1 | ||
fi | ||
|
||
exit 0 | ||
} | ||
|
||
clean() | ||
{ | ||
rm -f $STATUS_FILE | ||
rm -f $PORT_FILE | ||
rm -rf vine-run-info | ||
exit 0 | ||
} | ||
|
||
|
||
dispatch "$@" |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,60 @@ | ||
#! /usr/bin/env python | ||
|
||
import sys | ||
import ndcctools.taskvine as vine | ||
|
||
port_file = None | ||
try: | ||
port_file = sys.argv[1] | ||
except IndexError: | ||
sys.stderr.write("Usage: {} PORTFILE\n".format(sys.argv[0])) | ||
raise | ||
|
||
def main(): | ||
executor = vine.FuturesExecutor( | ||
port=[9123, 9129], manager_name="vine_hof_test", factory=False | ||
) | ||
|
||
print("listening on port {}".format(executor.manager.port)) | ||
with open(port_file, "w") as f: | ||
f.write(str(executor.manager.port)) | ||
|
||
nums = list(range(101)) | ||
|
||
rows = 3 | ||
mult_table = executor.allpairs(lambda x, y: x*y, range(rows), nums, chunk_size=11).result() | ||
assert sum(mult_table[1]) == sum(nums) | ||
assert sum(sum(r) for r in mult_table) == sum(sum(nums) * n for n in range(rows)) | ||
|
||
doubles = executor.map(lambda x: 2*x, nums, chunk_size=10).result() | ||
assert sum(doubles) == sum(nums)*2 | ||
|
||
doubles = executor.map(lambda x: 2*x, nums, chunk_size=13).result() | ||
assert sum(doubles) == sum(nums)*2 | ||
|
||
maximum = executor.reduce(max, nums, fn_arity=2).result() | ||
assert maximum == 100 | ||
|
||
maximum = executor.reduce(max, nums, fn_arity=25).result() | ||
assert maximum == 100 | ||
|
||
maximum = executor.reduce(max, nums, fn_arity=1000).result() | ||
assert maximum == 100 | ||
|
||
maximum = executor.reduce(max, nums, fn_arity=2, chunk_size=50).result() | ||
assert maximum == 100 | ||
|
||
minimum = executor.reduce(min, nums, fn_arity=2, chunk_size=50).result() | ||
assert minimum == 0 | ||
|
||
total = executor.reduce(sum, nums, fn_arity=11, chunk_size=13).result() | ||
assert total == sum(nums) | ||
|
||
|
||
|
||
|
||
if __name__ == "__main__": | ||
main() | ||
|
||
|
||
# vim: set sts=4 sw=4 ts=4 expandtab ft=python: |