google / google/xarray-beam

`OverflowError` when a chunk exceeds 2 GB

Open
#196 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Python
Stars
170
Forks
15
Avg merge
18h 27m
Merged PRs (30d)
1

Description

`DatasetCoder.estimate_size` returns `value.nbytes` directly ([`core.py#L275`](https://github.com/google/xarray-beam/blob/main/xarray_beam/_src/core.py#L275)). For
chunks larger than `2**31 - 1` bytes (~2 GB), this overflows the C `int` in Beam's Cython `CallbackCoderImpl.estimate_size` and raises:

```
OverflowError: value too large to convert to int
```

Noticed this hitting after upgrading to the latest version when loading datasets from GRIBs with multple atmostpheric levels, e.g., ECMWF IFS ensembles (51 members × 13 levels × 721 × 1440), where a single per-key `xarray.Dataset` chunk exceeds 2 GB before any rechunk.

Minimal repro:

```python
import apache_beam as beam
import numpy as np
import xarray as xr
import xarray_beam as xbeam

ds = xr.Dataset({"x": (("a", "b"), np.zeros((20000, 14000)))}) # ~2.1 GB

with beam.Pipeline() as p:
_ = (
p
| beam.Create([(xbeam.Key({}), ds)])
| beam.Reshuffle()
)
```

Capping the returned estimate at `2**31 - 1` (with a warning so oversized chunks are still visible) avoids the crash, will put up a PR.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.