jay
12/14/2024, 12:28 PMbatch_size affects peak memory usage for a simple UDF like so that returns a Python list of strings:
@daft.udf(return_dtype=str, batch_size=128)
def to_pylist_identity(s):
data = s.to_pylist()
# mock "expensive" functionality -- make a big array of strings
result = pa.compute.binary_repeat(pa.array(data), 8)
return data
• `batch_size=None`: Each Ray task uses ~290+MB at peak (bottlenecked by the pa.compute.binary_repeat(pa.array(data), 8) code)
• `batch_size=128`: Each Ray task still uses 100MB at peak (why??)
After some investigation:
• The batch_size parameter is indeed working as intended and the code is no longer bottlenecked at pa.compute.binary_repeat(pa.array(data), 8)
• HOWEVER it seems that the conversion of UDF pylist outputs to Series is now the bottleneck 😢 — Literally changing the above code to just do return pa.array(data) instead brings peak memory utilization of the Ray task to 8MB