huggingface / huggingface/datasets
`Dataset.map(num_proc=N)` worker crashes with `ValueError: I/O operation on closed file` when `finalize()` is interrupted
- Dominant language
- Python
- Stars
- 22k
- Forks
- 3.4k
- Avg merge
- 5d 7h
- Merged PRs (30d)
- 17
Description
### Describe the bug
We hit this while running Megatron-LM training through ms-swift, which uses `datasets` under the hood for preprocessing. Our preprocessing step calls `Dataset.map()` with a fairly aggressive `num_proc=512` on a ~25k-example JSONL file:
```bash
swift --dataset_num_proc 512 --dataset large.jsonl --max_length 131072 ...
```
Near the end of the map (around 94%), one of the worker processes died with a confusing two-stage traceback, and since this runs inside distributed training, the whole 128-GPU job went down with it.
#### Error logs
The worker first hit an `InterruptedError` while finalizing its output shard:
```
Map (num_proc=512): 94%|█████████▍| 24277/25742 [01:00<00:03, 398.48 examples/s]
[rank48]: Traceback (most recent call last):
[rank48]: File ".../datasets/arrow_dataset.py", line 4022, in _map_single
[rank48]: writer.finalize() # close_stream=bool(buf_writer is None))
[rank48]: File ".../datasets/arrow_writer.py", line 778, in finalize
[rank48]: self.stream.close()
[rank48]: File ".../fsspec/implementations/local.py", line 450, in close
[rank48]: return self.f.close()
[rank48]: InterruptedError: [Errno 4] Interrupted system call
```
That part is at least understandable — with 512 processes being spawned/reaped there's a lot of signal noise, and apparently `close()` can legitimately get interrupted.
The real problem is what happened next. The `except` block in `_map_single()` tried to call `finalize()` again, and that second call blew up with a completely different error:
```
[rank48]: During handling of the above exception, another exception occurred:
[rank48]:
[rank48]: Traceback (most recent call last):
[rank48]: File ".../multiprocess/pool.py", line 125, in worker
[rank48]: result = (True, func(*args, **kwds))
[rank48]: File ".../datasets/utils/py_utils.py", line 585, in _write_generator_to_queue
[rank48]: for i, result in enumerate(func(**kwargs)):
[rank48]: File ".../datasets/arrow_dataset.py", line 4027, in _map_single
[rank48]: writer.finalize()
[rank48]: File ".../datasets/arrow_writer.py", line 773, in finalize
[rank48]: self._build_writer(self.schema)
[rank48]: File ".../datasets/arrow_writer.py", line 601, in _build_writer
[rank48]: self.pa_writer = pa.RecordBatchStreamWriter(self.stream, self._schema)
[rank48]: File "pyarrow/ipc.py", line 90, in __init__
[rank48]: self._open(sink, schema, options=options)
[rank48]: ValueError: I/O operation on closed file
```
So the exception that actually kills the worker (`ValueError: I/O operation on closed file`) is not the original failure — it's produced by the error-handling path itself, which masks the real problem and makes this very confusing to debug.
The same script runs fine with smaller `num_proc` values (we've used 32 and 64 without issues), so this seems related to the high process count.
### Steps to reproduce the bug
```python
from datasets import load_dataset
ds = load_dataset("json", data_files="large.jsonl", split="train")
ds = ds.map(lambda x: x, num_proc=512)
```
It's timing/signal dependent so it doesn't fire 100% of the time, but with a high enough `num_proc` and enough examples it reproduces pretty reliably for us — we've hit it multiple times on ~25k examples with `num_proc=512`.
### Expected behavior
- A transient `InterruptedError` while closing the output stream shouldn't bring down the worker.
- The cleanup path in `_map_single()` shouldn't raise a *new* exception that hides the original one.
### Environment info
- `datasets`: 5.0.0 (I also checked the 5.0.1 source and the relevant code looks unchanged)
- Python: 3.12.11
- OS: Linux (training container)
- PyTorch: 2.8.0
- `multiprocess`: bundled with `datasets`
Contributor guide
Research direction
Start in datasets/arrow_dataset.py at _map_single() and inspect the exception path around writer.finalize(); then read datasets/arrow_writer.py finalize() and its stream handling. Reproduce with a large JSONL dataset and high num_proc, such as the provided 512-process example. Done means an interrupted close does not trigger a second closed-file exception that masks the original error.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- python
- Domain
- data, distributed-systems
- Issue type
- Bug
- Difficulty
- 3/5
- Estimated time
- 1-2 days
- Activity status
- Active
- Clarity
- Mostly clear
- Newbie friendliness
- 70/100