| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
feat: Channel-less intermediate op (#5999) ## Changes Made There is currently a bug with IntermediateOps where the worker that produces a HasMoreOutput is not being polled from again, in the maintain_order = True case. ### Recap Quick recap, what is actually going in with intermediate ops, and why does this only affect HasMoreOutput. <img width="1807" height="1578" alt="image" src="https://github.com/user-attachments/assets/b06d84b7-dd7b-443e-99c6-6dd14f05b8e4" /> We have intermediate op workers, that are responsible for actually executing the op on input morsels, a dispatcher that dispatches work to them, and a receiver that receives results. in the maintain order = True case, we use a round robin dispatcher that dispatches work in round robin fashion, and an OrderingAwareReceiver that receives results in a round robin fashion. ### Problem If a worker produces a result with HasMoreOutput, the round robin receiver SHOULD PULL FROM THAT WORKER AGAIN, BUT IT DOES NOT and instead it moves on to the next worker. This can create a deadlock where the worker that produces has more output wants to send again, but the round robin reciever is not awaiting it, and is instead awaiting from some other worker that could be, say, awaiting it's own input. ### Solution Just give the OrderingAwareReceiver info that it should pull from a certain receiver again. Bleh, this is now so messy. So, lets just get rid of the channels altogether. No more channels. We can remodel the intermediate op to a single concurrent state machine. while has_input or has_active_tasks { tokio::select { new_input = input_rx.recv() => active_tasks.spawn(execute(new_input)), result = active_tasks.join_next() => process_result(result) # if has more output just spawn again. } } This allows us to get rid channels within the intermediate op altogether, and only have a single channel connecting intermdiate ops. ### Results This script simulates a long pipeline of connecting intermediate ops. Theoretically less channels means less intermediate data. import daft import os @daft.func def generate_10_mb_data(x: int) -> bytes: return os.urandom(10 * 1024 * 1024) @daft.func def noop(x: bytes) -> bytes: return x df = daft.from_pydict({"id": [i for i in range(100)]}).into_batches(1) # Generate 10 MB of data for each row df = df.with_column("data_0", generate_10_mb_data(daft.col("id"))) # Add some noop udfs to add a lot of intermdiate ops for i in range(1, 10): df = df.with_column(f"data_{i}", noop(daft.col(f"data_{i-1}"))) # Send them to the void for p in df.iter_partitions(): pass Before: <img width="1113" height="367" alt="Screenshot 2026-01-09 at 1 37 21 PM" src="https://github.com/user-attachments/assets/2ed003fc-7764-4f32-8c2f-1ab571d1696a" /> After: <img width="1094" height="390" alt="Screenshot 2026-01-09 at 1 36 47 PM" src="https://github.com/user-attachments/assets/4faa9a91-218b-4c86-be00-983b67e20b98" /> We reduced peak memory by half. ### Future This HasMoreOutput bug affects streaming sinks as well actually. So if this fix works we should implement it on streaming sink as well, and then might as well do so for blocking sinks and we can get rid of the intra-op channels altogether. ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 8 个月前 | |
feat: Channel-less intermediate op (#5999) ## Changes Made There is currently a bug with IntermediateOps where the worker that produces a HasMoreOutput is not being polled from again, in the maintain_order = True case. ### Recap Quick recap, what is actually going in with intermediate ops, and why does this only affect HasMoreOutput. <img width="1807" height="1578" alt="image" src="https://github.com/user-attachments/assets/b06d84b7-dd7b-443e-99c6-6dd14f05b8e4" /> We have intermediate op workers, that are responsible for actually executing the op on input morsels, a dispatcher that dispatches work to them, and a receiver that receives results. in the maintain order = True case, we use a round robin dispatcher that dispatches work in round robin fashion, and an OrderingAwareReceiver that receives results in a round robin fashion. ### Problem If a worker produces a result with HasMoreOutput, the round robin receiver SHOULD PULL FROM THAT WORKER AGAIN, BUT IT DOES NOT and instead it moves on to the next worker. This can create a deadlock where the worker that produces has more output wants to send again, but the round robin reciever is not awaiting it, and is instead awaiting from some other worker that could be, say, awaiting it's own input. ### Solution Just give the OrderingAwareReceiver info that it should pull from a certain receiver again. Bleh, this is now so messy. So, lets just get rid of the channels altogether. No more channels. We can remodel the intermediate op to a single concurrent state machine. while has_input or has_active_tasks { tokio::select { new_input = input_rx.recv() => active_tasks.spawn(execute(new_input)), result = active_tasks.join_next() => process_result(result) # if has more output just spawn again. } } This allows us to get rid channels within the intermediate op altogether, and only have a single channel connecting intermdiate ops. ### Results This script simulates a long pipeline of connecting intermediate ops. Theoretically less channels means less intermediate data. import daft import os @daft.func def generate_10_mb_data(x: int) -> bytes: return os.urandom(10 * 1024 * 1024) @daft.func def noop(x: bytes) -> bytes: return x df = daft.from_pydict({"id": [i for i in range(100)]}).into_batches(1) # Generate 10 MB of data for each row df = df.with_column("data_0", generate_10_mb_data(daft.col("id"))) # Add some noop udfs to add a lot of intermdiate ops for i in range(1, 10): df = df.with_column(f"data_{i}", noop(daft.col(f"data_{i-1}"))) # Send them to the void for p in df.iter_partitions(): pass Before: <img width="1113" height="367" alt="Screenshot 2026-01-09 at 1 37 21 PM" src="https://github.com/user-attachments/assets/2ed003fc-7764-4f32-8c2f-1ab571d1696a" /> After: <img width="1094" height="390" alt="Screenshot 2026-01-09 at 1 36 47 PM" src="https://github.com/user-attachments/assets/4faa9a91-218b-4c86-be00-983b67e20b98" /> We reduced peak memory by half. ### Future This HasMoreOutput bug affects streaming sinks as well actually. So if this fix works we should implement it on streaming sink as well, and then might as well do so for blocking sinks and we can get rid of the intra-op channels altogether. ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 8 个月前 |
| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
| 8 个月前 | ||
| 8 个月前 |