# Prepare ESA CCI biomass data

In [None]:
# Libraries
import os, time, shutil, rioxarray
import numpy as np
import xarray as xr
from scipy import ndimage
import matplotlib.pyplot as plt

from dask.distributed import Lock

In [None]:
# Directories
dir_data =  '../data/'
dir01 = '../paper_deficit/output/01_prep/'

---

### Pre-processing

In [None]:
# Libraries
from dask_jobqueue import SLURMCluster
from dask.distributed import Client
import dask

# Initialize dask
cluster = SLURMCluster(
    queue='compute',                      # SLURM queue to use
    cores=24,                             # Number of CPU cores per job
    memory='256 GB',                      # Memory per job
    account='bm0891',                     # Account allocation
    interface="ib0",                      # Network interface for communication
    walltime='01:00:00',                  # Maximum runtime per job
    local_directory='../dask/',           # Directory for local storage
    job_extra_directives=[                # Additional SLURM directives for logging
        '-o ../dask/LOG_worker_%j.o',     # Output log
        '-e ../dask/LOG_worker_%j.e'      # Error log
    ]
)

# Scale dask cluster
cluster.scale(jobs=10)

# Configurate dashboard url
dask.config.config.get('distributed').get('dashboard').update(
    {'link': '{JUPYTERHUB_SERVICE_PREFIX}/proxy/{port}/status'}
)

# Create client
client = Client(cluster)

client

In [None]:
# Function to trim memory of workers
def trim_memory() -> int:
    import ctypes
    libc = ctypes.CDLL("libc.so.6")
    return libc.malloc_trim(0)

In [None]:
def prep_esabio():

    """
    Prepare ESA CCI biomass data for regridding
    """
    
    for var in ['agb', 'agbsd']:
        for year in [2015, 2016, 2017, 2018, 2019, 2020, 2021]:
            # File path
            file_path = dir_data + 'esa_biomass/v5_01/' + \
                'ESACCI-BIOMASS-L4-AGB-MERGED-1000m-fv5.01.nc'
            # Get data
            ds = xr.open_dataset(file_path).chunk(dict(lat=5000, lon=5000))
            # Select array
            if var == 'agb':
                da = ds['agb'].rename('esabio_agb' + str(year))
            if var == 'agbsd':
                da = ds['agb_sd'].rename('esabio_agbsd' + str(year))
            # Select year and rename
            da = da.sel(time=str(year)).squeeze('time').drop_vars('time').astype(np.float32)
            # Remove attributes (creates problems with CDO)
            da.attrs = {}
            # Export as tif
            da.rio.to_raster(dir01 + 'esabio_' + var + str(year) + '.tif')


%time prep_esabio()

In [None]:
# Close dask cluster
cluster.close()

In [None]:
# Wait for 60s for dask client to completely disconnect
time.sleep(60)

---

### Regridding

In [None]:
# Import regridding function
from regrid_high_res_v1_01 import regrid_high_res, prep_tif

In [None]:
def regrid_da(f_source, dir_target, dir_source, dir_out, 
              size_tiles, fill_value=None, olap=1):  
    """Regrid large xarray dataarrays.

    Args:
        f_source (str): The filename (without extension) of the source .tif file to be regridded.
        dir_target (str): Directory containing target grid .tif file.
        dir_source (str): Directory containing the the source  .tif file.
        dir_out (str): Directory to store the output and intermediate files.
        size_tiles (int): Size of the regridding tiles in degrees.
        fill_value (float, optional): Fill value to use in the regridding process. Defaults to None.
        olap (int, optional): Overlap size in degrees for regridding tiles. Defaults to 1.
        
    Returns:
        xarray.Dataset: The combined dataset after regridding.
    """
    # Prepare the target and source data arrays from TIFF files
    da_target = prep_tif(dir_target + 'target_grid.tif', 'target_grid')
    da_source = prep_tif(dir_source + f_source + '.tif', f_source)
    # Regridd source array to target grid
    regrid_high_res(da_target, da_source, dir_out,
                    account='bm0891', partition='compute',
                    size_tiles=size_tiles, olap=olap, fill_value = fill_value,
                    type_export='zarr', del_interm=False)

In [None]:
for i in [2015, 2016, 2017, 2018, 2019, 2020, 2021]:
    %time regrid_da('esabio_agb' + str(i), dir01, dir01, dir01, \
                    30, np.nan, 0.1)
# Takes about 6 min for one file

In [None]:
for i in [2015, 2016, 2017, 2018, 2019, 2020, 2021]:
    %time regrid_da('esabio_agbsd' + str(i), dir01, dir01, dir01, \
                    30, np.nan, 0.1)

---

### Fill nans

In [None]:
def fill_nans(var, dir_out):
    """
    Fills NaN values in the specified variable's dataset using the nearest valid 
    data, applies a land mask, and exports the result to a new Zarr dataset.

    Args:
        var (str): Name of the variable to process (e.g., 'temperature', 'precipitation').
        dir_out (str): Directory where prepared data is stored and the filled dataset will be exported.

    Returns:
        None
    """

    def fill_nans_array(data, invalid):
        """
        Replace invalid (NaN) data cells by the value of the nearest valid data 
        cell.
        """
        ind = ndimage.distance_transform_edt(invalid,
                                             return_distances=False,
                                             return_indices=True)
        return data[tuple(ind)]

    # Paths for input and output
    land_mask_path = os.path.join(dir_out, 'ds_prep_copernicus_land_mask.zarr')
    var_data_path = os.path.join(dir_out, f'ds_regridded_{var}.zarr')
    output_path = os.path.join(dir_out, f'ds_prep_{var}.zarr')

    # Read land mask data
    da_land = xr.open_zarr(land_mask_path) \
                .chunk(dict(lat=5000, lon=5000)) \
                .copernicus_land_mask \
                .compute()

    # Read variable data
    da_var = xr.open_zarr(var_data_path)['regridded_' + var]

    # Fill nan using function fill_nans_array
    # If there are no NaNs, skip filling process
    if not da_var.isnull().any():
        da_fill = da_var.values  # No filling required
    else:
        da_fill = fill_nans_array(da_var.values, da_var.isnull().values)

    # Create a new Dataset with filled data
    ds_filled = xr.Dataset(dict(lat = da_var.lat, lon=da_var.lon))
    ds_filled[var] = (('lat', 'lon'), da_fill)

    # Apply land mask to the filled data
    ds_filled = ds_filled.where(da_land)

    # Export the filled dataset to Zarr format
    ds_filled.chunk(dict(lat=5000, lon=5000)) \
             .to_zarr(output_path, mode='w')

In [None]:
for i in ['esabio_agb2015', 'esabio_agb2016', 'esabio_agb2017',
          'esabio_agb2018', 'esabio_agb2019', 'esabio_agb2020',
          'esabio_agb2021']:
    %time fill_nans(i, dir01)

In [None]:
for i in ['esabio_agbsd2015', 'esabio_agbsd2016', 'esabio_agbsd2017',
          'esabio_agbsd2018', 'esabio_agbsd2019', 'esabio_agbsd2020',
          'esabio_agbsd2021']:
    %time fill_nans(i, dir01)

---

### Check

In [None]:
# Plot to check
for i in ['esabio_agb2015', 'esabio_agb2016', 'esabio_agb2017', 
          'esabio_agb2018', 'esabio_agb2019', 'esabio_agb2020',
          'esabio_agb2021']:
    fig, ax = plt.subplots(figsize=(10, 5), ncols=1, nrows=1)
    xr.open_zarr(dir01 + 'ds_prep_' + i + '.zarr')[i] \
        .plot.imshow(ax=ax, robust=True)
    ax.set_title(i)

In [None]:
# Plot to check
for i in ['esabio_agbsd2015', 'esabio_agbsd2016', 'esabio_agbsd2017', 
          'esabio_agbsd2018', 'esabio_agbsd2019', 'esabio_agbsd2020',
          'esabio_agbsd2021']:
    fig, ax = plt.subplots(figsize=(10, 5), ncols=1, nrows=1)
    xr.open_zarr(dir01 + 'ds_prep_' + i + '.zarr')[i] \
        .plot.imshow(ax=ax, robust=True)
    ax.set_title(i)