Xarray multiindex verilerini parçalar halinde yazma

Sep 15 2020

Büyük bir çok boyutlu veri kümesini verimli bir şekilde yeniden yapılandırmaya çalışıyorum. Piksel konumu için xy koordinatlarına, görüntü alma zamanı için zamana ve toplanan farklı veriler için bantlara sahip bir dizi bant ile zaman içinde uzaktan algılanan bir dizi resme sahip olduğumu varsayalım.

Benim kullanım durumumda, xarray koordinat uzunluklarının yaklaşık olarak x (3000), y (3000), zaman (10) ve kayan nokta verisi bantları (40) olduğunu varsayalım. Yani 100 gb + veri.

Bu örnekten çalışmaya çalışıyorum ama onu bu davaya çevirmekte sorun yaşıyorum.

Küçük veri kümesi örneği

NOT: gerçek veriler bu örnekten çok daha büyüktür.

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'

Yeniden yapılandırılmış - istiflenmiş ve yeri değiştirilmiş

Aşağıdakilerin çıktılarını saklamam gerekiyor:

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'

Ben DASK kullanmak umuduyla ve erişilebilir büyük kümeler halinde diske sonucu yazma xarray am open_mfdataset . parke iyi bir seçenek gibi görünüyor, ancak onu parçalar halinde nasıl yazacağımı çözemiyorum (src, bellekte saklamak için çok büyük).

@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)

Bu noktada sadece bazı iyi fikirler umuyorum. Teşekkürler!

Yanıtlar

3 dcherian Sep 30 2020 at 01:26

Burada bir çözümüm var (https://github.com/pydata/xarray/issues/1077#issuecomment-644803374) çoklu dizine alınmış veri kümelerini dosyaya yazmak için.

Veri kümesini netCDF olarak yazılabilen bir forma manuel olarak "kodlamanız" gerekir. Ve sonra tekrar okuduğunuzda "kodunu çöz".

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

Şimdilik çözüm bu değil, kodunuzun bir sürümü, başkaları bu sorunu çözmeye çalışırsa kolayca yeniden üretilebilecek şekilde değiştirildi:

Sorun, stack( concatenated.stack(sample=('y','x','time')) işlemi ile ilgilidir . Bu adımda hafıza artmaya devam ediyor ve süreç artıyor killed.

concatenatedNesne "Dask destekli" dir xarray.DataArray. Yani stackoperasyonun Dask tarafından tembel bir şekilde yapılmasını bekleyebiliriz . Öyleyse, süreç neden killedbu adımda?

Burada olup bitenler için 2 olasılık:

  • İşlem stackaslında Dask tarafından tembelce yapılır, ancak veriler çok büyük olduğu için, Dask için gereken minimum bellek bile çok fazladır.

  • İşlem stackDask destekli DEĞİLDİR


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

Bir değerlerini değiştirmek deneyebilirsiniz nrowsve ncolsboyutunu değiştirmek amacıyla concatenated. Ve performans için chunksde değişebiliriz / değiştirmeliyiz .

Not: Bunu bile denedim

concatenated.to_netcdf("concatenated.nc")
concatenated = xr.open_dataarray("concatenated.nc", chunks=10)

Bu, Dask destekli bir DataArray olduğundan emin olmak ve parçaları da ayarlayabilmek içindir. Şunun için farklı değer / ler denedim chunksama her zaman bellek yetersiz.