Lightning-AI / Lightning-AI/litData

Using a streaming dataloader with an unbalanced dataset yields unexpected batch sizes.

Open
#199 10 comments 2 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

bug help wanted
Dominant language
Python
Stars
614
Forks
106
Avg merge
15h 8m
Merged PRs (30d)
22

Description

## 🐛 Bug

I have two datasets which are unbalanced, where one dataset is 1000x larger than the other. I would like to sample from two of the datasets such that the ratio of samples from each is 1:100. When doing so, the batches are of irregular size are returned during iteration.

I think there are 2 issues which this test surfaces:
1) The first batch returned by each worker is not properly sized.
2) `drop_last` does not appear to work as intended, since the last batch is not a full sized batch

I don't think this is related to #179, but it's possible

I've been attempting to fix this, but I'm not sure what the root of the issue is. I would be very appreciative if you could fix this or point me in the right direction.

Thanks!

### To Reproduce

```
@pytest.mark.skipif(sys.platform == "win32", reason="too slow in CI")
def test_unbalanced_combined_dataset_with_dataloader(tmpdir):
data_dir_1 = os.path.join(tmpdir, "data_1")
data_dir_2 = os.path.join(tmpdir, "data_2")
cache_dir_1 = os.path.join(tmpdir, "cache_dir_1")
cache_dir_2 = os.path.join(tmpdir, "cache_dir_2")

os.makedirs(data_dir_1)
os.makedirs(data_dir_2)
os.makedirs(cache_dir_1)
os.makedirs(cache_dir_2)

cache = Cache(input_dir=str(data_dir_1), chunk_size=2)

for i in range(10):
cache[i] = i

cache.done()
cache.merge()

cache = Cache(input_dir=str(data_dir_2), chunk_size=2)

for i in range(10000):
cache[i] = i + 10

cache.done()
cache.merge()

dataset1 = StreamingDataset(input_dir=Dir(cache_dir_1, data_dir_1), shuffle=True)
dataset2 = StreamingDataset(input_dir=Dir(cache_dir_2, data_dir_2), shuffle=True)
dataset = CombinedStreamingDataset(
datasets=[dataset1, dataset2], weights=[0.01, 0.99], iterate_over_all=False, seed=12345
)
dataloader = StreamingDataLoader(dataset, num_workers=3, batch_size=100, drop_last=True, persistent_workers=True, shuffle=True, prefetch_factor=2)

assert dataset1.current_epoch == 1
assert dataset2.current_epoch == 1

batches_1 = []
batch_sizes_1 = []
for batch in dataloader:
batch_sizes_1.append(batch.size(0))
batches_1.append(batch)

assert batch_sizes_1[2] == 91
assert batch_sizes_1[-1] == 40
# This will fail since the third and last index are not 100. (Above 2 assertions pass)
assert batch_sizes_1 == [100 for _ in batches_1]

```

### Expected behavior

All batch sizes should be the same.

### Additional context

This issue is independent of whether `drop_last`, `shuffle`, and `persistent_workers` are set to True or False

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start by running test_unbalanced_combined_dataset_with_dataloader, then inspect StreamingDataLoader and CombinedStreamingDataset with the multi-worker settings in the reproduction. Trace how worker batches are assembled and how drop_last is applied. Done means the reproduced iteration returns only batches of size 100, including the first and final batches.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
data
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.