xarray 다중 인덱스 데이터를 청크로 쓰기
대규모 다차원 데이터 세트를 효율적으로 재구성하려고합니다. 픽셀 위치에 대한 좌표 xy, 이미지 획득 시간에 대한 시간 및 수집 된 서로 다른 데이터에 대한 밴드를 사용하여 시간이 지남에 따라 원격으로 감지 된 여러 이미지가 있다고 가정 해 보겠습니다.
내 사용 사례에서는 xarray 좌표 길이가 대략 x (3000), y (3000), 시간 (10)이고 부동 소수점 데이터 밴드 (40)가 있다고 가정합니다. 따라서 100GB 이상의 데이터.
이 예제 에서 작업하려고 했지만이 사례로 번역하는 데 문제가 있습니다.
작은 데이터 세트 예
참고 : 실제 데이터는이 예보다 훨씬 큽니다.
import numpy as np
import dask.array as da
import xarray as xr
nrows = 100
ncols = 200
row_chunks = 50
col_chunks = 50
data = da.random.random(size=(1, nrows, ncols), chunks=(1, row_chunks, col_chunks))
def create_band(data, x, y, band_name):
return xr.DataArray(data,
dims=('band', 'y', 'x'),
coords={'band': [band_name],
'y': y,
'x': x})
def create_coords(data, left, top, celly, cellx):
nrows = data.shape[-2]
ncols = data.shape[-1]
right = left + cellx*ncols
bottom = top - celly*nrows
x = np.linspace(left, right, ncols) + cellx/2.0
y = np.linspace(top, bottom, nrows) - celly/2.0
return x, y
x, y = create_coords(data, 1000, 2000, 30, 30)
src = []
for time in ['t1', 't2', 't3']:
src_t = xr.concat([create_band(data, x, y, band) for band in ['blue', 'green', 'red', 'nir']], dim='band')\
.expand_dims(dim='time')\
.assign_coords({'time': [time]})
src.append(src_t)
src = xr.concat(src, dim='time')
print(src)
<xarray.DataArray 'random_sample-5840d8564d778d573dd403f27c3f47a5' (time: 3, band: 4, y: 100, x: 200)>
dask.array<concatenate, shape=(3, 4, 100, 200), dtype=float64, chunksize=(1, 1, 50, 50), chunktype=numpy.ndarray>
Coordinates:
* x (x) float64 1.015e+03 1.045e+03 1.075e+03 ... 6.985e+03 7.015e+03
* band (band) object 'blue' 'green' 'red' 'nir'
* y (y) float64 1.985e+03 1.955e+03 1.924e+03 ... -984.7 -1.015e+03
* time (time) object 't1' 't2' 't3'
재구성-스택 및 전치
다음의 출력을 저장해야합니다.
print(src.stack(sample=('y','x','time')).T)
<xarray.DataArray 'random_sample-5840d8564d778d573dd403f27c3f47a5' (sample: 60000, band: 4)>
dask.array<transpose, shape=(60000, 4), dtype=float64, chunksize=(3600, 1), chunktype=numpy.ndarray>
Coordinates:
* band (band) object 'blue' 'green' 'red' 'nir'
* sample (sample) MultiIndex
- y (sample) float64 1.985e+03 1.985e+03 ... -1.015e+03 -1.015e+03
- x (sample) float64 1.015e+03 1.015e+03 ... 7.015e+03 7.015e+03
- time (sample) object 't1' 't2' 't3' 't1' 't2' ... 't3' 't1' 't2' 't3'
dask와 xarray를 사용하여 결과를 디스크에 청크 단위로 쓰고 open_mfdataset에 액세스 할 수 있기 를 바랍니다 . parquet은 좋은 옵션처럼 보이지만 청크로 작성하는 방법을 알 수 없습니다 (src가 메모리에 저장하기에 너무 큽니다).
@dask.delayed
def stacker(data):
return data.stack(sample=('y','x','time')).T.to_pandas()
stacker(src).to_parquet('out_*.parquet')
def stack_write(data):
data.stack(sample=('y','x','time')).T.to_pandas().to_parquet('out_*.parquet')
return None
stack_write(src)
이 시점에서 나는 단지 좋은 아이디어를 바라고 있습니다. 감사!
답변
여기에 해결책이 있습니다 (https://github.com/pydata/xarray/issues/1077#issuecomment-644803374) 파일에 다중 인덱싱 된 데이터 세트를 작성합니다.
netCDF로 작성할 수있는 양식으로 데이터 세트를 수동으로 "인코딩"해야합니다. 그리고 다시 읽을 때 "디코딩"합니다.
import numpy as np
import pandas as pd
import xarray as xr
def encode_multiindex(ds, idxname):
encoded = ds.reset_index(idxname)
coords = dict(zip(ds.indexes[idxname].names, ds.indexes[idxname].levels))
for coord in coords:
encoded[coord] = coords[coord].values
shape = [encoded.sizes[coord] for coord in coords]
encoded[idxname] = np.ravel_multi_index(ds.indexes[idxname].codes, shape)
encoded[idxname].attrs["compress"] = " ".join(ds.indexes[idxname].names)
return encoded
def decode_to_multiindex(encoded, idxname):
names = encoded[idxname].attrs["compress"].split(" ")
shape = [encoded.sizes[dim] for dim in names]
indices = np.unravel_index(encoded.landpoint.values, shape)
arrays = [encoded[dim].values[index] for dim, index in zip(names, indices)]
mindex = pd.MultiIndex.from_arrays(arrays)
decoded = xr.Dataset({}, {idxname: mindex})
for varname in encoded.data_vars:
if idxname in encoded[varname].dims:
decoded[varname] = (idxname, encoded[varname].values)
return decoded
이것은 현재 해결책이 아니라 다른 사람들이이 문제를 해결하려고 할 때 쉽게 재현 할 수 있도록 수정 된 코드 버전입니다.
stack작업 ( concatenated.stack(sample=('y','x','time'))에 문제가 있습니다. 이 단계에서 메모리는 계속 증가하고 프로세스는 killed입니다.
concatenated객체는 "DASK는 백업"입니다 xarray.DataArray. 따라서 stack작업이 Dask에 의해 느리게 수행 될 것으로 예상 할 수 있습니다. 그렇다면 killed이 단계 의 프로세스 는 무엇입니까?
여기서 일어나는 일에 대한 두 가지 가능성 :
stack작업은 사실 DASK에 의해 느리게 수행되지만 데이터가 매우 때문에 큰는 DASK 심지어 필요한 최소 메모리는 너무 많은 것을이
stack작업은 Dask 지원이 아닙니다.
import numpy as np
import dask.array as da
import xarray as xr
from numpy.random import RandomState
nrows = 20000
ncols = 20000
row_chunks = 500
col_chunks = 500
# Create a reproducible random numpy array
prng = RandomState(1234567890)
numpy_array = prng.rand(1, nrows, ncols)
data = da.from_array(numpy_array, chunks=(1, row_chunks, col_chunks))
def create_band(data, x, y, band_name):
return xr.DataArray(data,
dims=('band', 'y', 'x'),
coords={'band': [band_name],
'y': y,
'x': x})
def create_coords(data, left, top, celly, cellx):
nrows = data.shape[-2]
ncols = data.shape[-1]
right = left + cellx*ncols
bottom = top - celly*nrows
x = np.linspace(left, right, ncols) + cellx/2.0
y = np.linspace(top, bottom, nrows) - celly/2.0
return x, y
x, y = create_coords(data, 1000, 2000, 30, 30)
bands = ['blue', 'green', 'red', 'nir']
times = ['t1', 't2', 't3']
bands_list = [create_band(data, x, y, band) for band in bands]
src = []
for time in times:
src_t = xr.concat(bands_list, dim='band')\
.expand_dims(dim='time')\
.assign_coords({'time': [time]})
src.append(src_t)
concatenated = xr.concat(src, dim='time')
print(concatenated)
# computed = concatenated.compute() # "computed" is ~35.8GB
stacked = concatenated.stack(sample=('y','x','time'))
transposed = stacked.T
의 크기를 변경하기 위해 nrows및 의 값을 변경하려고 할 수 있습니다 . 그리고 성능을 위해 우리는 너무 다양 할 수 / 있어야합니다 .ncolsconcatenatedchunks
참고 : 나는 이것을 시도했다
concatenated.to_netcdf("concatenated.nc")
concatenated = xr.open_dataarray("concatenated.nc", chunks=10)
이는 Dask 지원 DataArray인지 확인하고 청크도 조정할 수 있도록하기위한 것입니다. 에 대해 다른 값을 시도 chunks했지만 항상 메모리가 부족합니다.