xarray 다중 인덱스 데이터를 청크로 쓰기

Sep 15 2020

대규모 다차원 데이터 세트를 효율적으로 재구성하려고합니다. 픽셀 위치에 대한 좌표 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)

이 시점에서 나는 단지 좋은 아이디어를 바라고 있습니다. 감사!

답변

3 dcherian Sep 30 2020 at 01:26

여기에 해결책이 있습니다 (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
1 Rivers Nov 15 2020 at 18:20

이것은 현재 해결책이 아니라 다른 사람들이이 문제를 해결하려고 할 때 쉽게 재현 할 수 있도록 수정 된 코드 버전입니다.

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했지만 항상 메모리가 부족합니다.