Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 29 additions & 19 deletions docs/guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,11 +67,12 @@ types. Configuration sequences are again merged into one.
import os
from zappend.api import zappend

zappend(os.listdir("inputs"),
config=["configs/base.yaml",
"configs/mycube.yaml"],
target_dir="outputs/mycube.zarr",
dry_run=True)
zappend(
os.listdir("inputs"),
config=["configs/base.yaml", "configs/mycube.yaml"],
target_dir="outputs/mycube.zarr",
dry_run=True,
)
```

The remainder of this guide explains the how to use the various `zappend`
Expand Down Expand Up @@ -463,8 +464,8 @@ You can compute `scale_factor` and `add_offset` from given data range in physica
according to

```python
add_offset = memory_value_min
scale_factor = (memory_value_max - memory_value_min) / (2 ** num_bits - 1)
add_offset = memory_value_min
scale_factor = (memory_value_max - memory_value_min) / (2**num_bits - 1)
```

with `num_bits` being the number of bits for the integer type to be used.
Expand Down Expand Up @@ -706,6 +707,7 @@ an `xarray.Dataset`:
```python
import xarray as xr


# Slice source argument `path` is just an example.
def get_dataset(path: str) -> xr.Dataset:
# Provide dataset here. No matter how, e.g.:
Expand All @@ -721,6 +723,7 @@ you can turn your slice source function into a
from contextlib import contextmanager
import xarray as xr


# Slice source argument `path` is just an example.
@contextmanager
def get_dataset(path: str) -> xr.Dataset:
Expand All @@ -747,6 +750,7 @@ You can also implement your slice source as a class derived from the abstract
import xarray as xr
from zappend.api import SliceSource


class MySliceSource(SliceSource):
# Slice source argument `path` is just an example.
def __init__(self, path: str):
Expand Down Expand Up @@ -781,9 +785,11 @@ qualified name of the slice source function or class:
If you use the `zappend` function, you can pass the function or class directly:

```python
zappend(["slice-1.nc", "slice-2.nc", "slice-3.nc"],
target_dir="target.zarr",
slice_source=MySliceSource)
zappend(
["slice-1.nc", "slice-2.nc", "slice-3.nc"],
target_dir="target.zarr",
slice_source=MySliceSource,
)
```

If the slice source setting is used, each slice item passed to `zappend` is passed as
Expand Down Expand Up @@ -835,6 +841,7 @@ from zappend.api import Context
from zappend.api import SliceSource
from zappend.api import zappend


class MySliceSource(SliceSource):
def __init__(self, ctx: Context, slice_path: str):
self.quantiles = ctx.config.extra.get("quantiles", [0.5])
Expand All @@ -849,25 +856,28 @@ class MySliceSource(SliceSource):
if self.ds is not None:
self.ds.close()

def get_agg_slice(self, slice_ds: xr.Dataset) -> xr.Dataset:
def get_agg_slice(self, slice_ds: xr.Dataset) -> xr.Dataset:
agg_slice_ds = slice_ds.quantile(self.quantiles, dim="time")
# Re-introduce time dimension of size one
agg_slice_ds = agg_slice_ds.expand_dims("time", axis=0)
agg_slice_ds.coords["time"] = self.get_mean_time(slice_ds)
return agg_slice_ds
return agg_slice_ds

@classmethod
def get_mean_time(cls, slice_ds: xr.Dataset) -> xr.DataArray:
time = slice_ds.time
t0 = time[0]
dt = time[-1] - t0
return xr.DataArray(np.array([t0 + dt / 2],
dtype=slice_ds.time.dtype),
dims="time")

zappend(["slice-1.nc", "slice-2.nc", "slice-3.nc"],
target_dir="target.zarr",
slice_source=MySliceSource)
return xr.DataArray(
np.array([t0 + dt / 2], dtype=slice_ds.time.dtype), dims="time"
)


zappend(
["slice-1.nc", "slice-2.nc", "slice-3.nc"],
target_dir="target.zarr",
slice_source=MySliceSource,
)
```

## Profiling
Expand Down
2 changes: 1 addition & 1 deletion docs/hooks/download_fonts.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# Copyright © 2024-2026 Brockmann Consult and contributors
# Copyright © 2024-2026 Brockmann Consult and contributors
# Permissions are hereby granted under the terms of the MIT License:
# https://opensource.org/licenses/MIT.

Expand Down
12 changes: 8 additions & 4 deletions docs/howdoi.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,20 +16,24 @@ import rioxarray as rxr
import xarray as xr
from zappend.api import zappend


def get_dataset_from_geotiff(tiff_path):
ds = rxr.open_rasterio(tiff_path)
# Add missing time dimension
slice_time = get_slice_time(tiff_path)
slice_time = get_slice_time(tiff_path)
slice_ds = ds.expand_dims("time", axis=0)
slice_ds.coords["time"] = xr.Dataset(np.array([slice_time]), dims="time")
try:
yield slice_ds
finally:
ds.close()

zappend(sorted(glob.glob("inputs/*.tif")),
slice_source=get_dataset_from_geotiff,
target_dir="output/tif-cube.zarr")

zappend(
sorted(glob.glob("inputs/*.tif")),
slice_source=get_dataset_from_geotiff,
target_dir="output/tif-cube.zarr",
)
```

In the example above, function `get_slice_time()` returns the time label
Expand Down
27 changes: 17 additions & 10 deletions docs/start.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,17 +66,21 @@ Process list of slices stored in S3 [configuration](config.md) in `config`:
```python
from zappend.api import zappend

config = {
config = {
"target_dir": "target.zarr",
"slice_storage_options": {
"key": "...",
"secret": "...",
}
"key": "...",
"secret": "...",
},
}

zappend((f"s3:/mybucket/data/{name}"
for name in ["slice-1.nc", "slice-2.nc", "slice-3.nc"]),
config=config)
zappend(
(
f"s3:/mybucket/data/{name}"
for name in ["slice-1.nc", "slice-2.nc", "slice-3.nc"]
),
config=config,
)
```

Slice items can also be arguments passed to your custom _slice source_,
Expand All @@ -91,9 +95,12 @@ def get_dataset(path: str):
ds = xr.open_dataset(path)
return ds.drop_vars(["ndvi_min", "ndvi_max"])

zappend(["slice-1.nc", "slice-2.nc", "slice-3.nc"],
slice_source=get_dataset,
target_dir="target.zarr")

zappend(
["slice-1.nc", "slice-2.nc", "slice-3.nc"],
slice_source=get_dataset,
target_dir="target.zarr",
)
```

For the details, please refer to the section [_Slice Sources_](guide.md#slice-sources) in the
Expand Down
77 changes: 39 additions & 38 deletions examples/zappend-demo.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -593,7 +593,9 @@
}
],
"source": [
"slice_ds.SPM.isel(time=0).sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(vmax=1.5)"
"slice_ds.SPM.isel(time=0).sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(\n",
" vmax=1.5\n",
")"
]
},
{
Expand Down Expand Up @@ -1768,7 +1770,9 @@
}
],
"source": [
"target_ds.SPM.mean(\"time\").sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(vmax=1.5)"
"target_ds.SPM.mean(\"time\").sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(\n",
" vmax=1.5\n",
")"
]
},
{
Expand Down Expand Up @@ -1799,7 +1803,9 @@
}
],
"source": [
"target_ds.SPM_QI.max(\"time\").sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(vmax=100.0)"
"target_ds.SPM_QI.max(\"time\").sel(lon=slice(9.2, 12), lat=slice(58.7, 56)).plot.imshow(\n",
" vmax=100.0\n",
")"
]
},
{
Expand Down Expand Up @@ -2959,7 +2965,7 @@
"outputs": [],
"source": [
"# Uncomment to show configuration refernce\n",
"#Markdown(get_config_schema(format=\"md\"))"
"# Markdown(get_config_schema(format=\"md\"))"
]
},
{
Expand All @@ -2970,44 +2976,31 @@
"outputs": [],
"source": [
"config = {\n",
" \"target_dir\": target_path, \n",
" \"target_dir\": target_path,\n",
" \"variables\": {\n",
" # We want the time coordinate variable to use a larger chunk size\n",
" # than the default (= 1 here)\n",
" \"time\": {\n",
" \"encoding\": {\n",
" \"chunks\": [100]\n",
" }\n",
" }\n",
" \"time\": {\"encoding\": {\"chunks\": [100]}}\n",
" },\n",
" # Log to the console.\n",
" # Note you could also configure the log output for dask here.\n",
" \"logging\": {\n",
" \"version\": 1,\n",
" \"formatters\": {\n",
" \"normal\": {\n",
" \"format\": \"%(asctime)s %(levelname)s %(message)s\",\n",
" \"style\": \"%\"\n",
" }\n",
" \"normal\": {\"format\": \"%(asctime)s %(levelname)s %(message)s\", \"style\": \"%\"}\n",
" },\n",
" \"handlers\": {\n",
" \"console\": {\n",
" \"class\": \"logging.StreamHandler\",\n",
" \"formatter\": \"normal\"\n",
" }\n",
" \"console\": {\"class\": \"logging.StreamHandler\", \"formatter\": \"normal\"}\n",
" },\n",
" \"loggers\": {\n",
" \"zappend\": {\n",
" \"level\": \"INFO\",\n",
" \"handlers\": [\"console\"]\n",
" }, \n",
" \"zappend\": {\"level\": \"INFO\", \"handlers\": [\"console\"]},\n",
" \"notebook\": {\n",
" # Will use this one later below\n",
" \"level\": \"INFO\",\n",
" \"handlers\": [\"console\"]\n",
" }\n",
" }\n",
" }\n",
" \"handlers\": [\"console\"],\n",
" },\n",
" },\n",
" },\n",
"}"
]
},
Expand All @@ -3029,6 +3022,7 @@
"outputs": [],
"source": [
"import yaml\n",
"\n",
"with open(config_path, mode=\"w\") as f:\n",
" yaml.dump(config, f)"
]
Expand Down Expand Up @@ -4189,6 +4183,7 @@
"outputs": [],
"source": [
"from logging import getLogger\n",
"\n",
"LOG = getLogger(\"notebook\")"
]
},
Expand All @@ -4209,11 +4204,13 @@
"source": [
"log = getLogger(\"notebook\")\n",
"\n",
"\n",
"def process_slice(slice_path: str) -> xr.Dataset:\n",
" LOG.info(f\"Processing slice {slice_path}\")\n",
" slice_ds = xr.open_dataset(slice_path) \n",
" slice_ds = xr.open_dataset(slice_path)\n",
" return slice_ds.drop_vars([\"SPM_QI\", \"TUR_QI\", \"crs\"])\n",
"\n",
"\n",
"slice_generator = (process_slice(slice_path) for slice_path in slice_paths)"
]
},
Expand Down Expand Up @@ -4999,16 +4996,17 @@
"source": [
"from zappend.api import SliceSource\n",
"\n",
"\n",
"class MySliceSource(SliceSource):\n",
" def __init__(self, slice_path):\n",
" self.slice_path = slice_path\n",
" self.slice_ds = None\n",
" \n",
" def get_dataset(self) -> xr.Dataset: \n",
"\n",
" def get_dataset(self) -> xr.Dataset:\n",
" LOG.info(f\"Processing slice {self.slice_path}\")\n",
" self.slice_ds = xr.open_dataset(self.slice_path)\n",
" return self.slice_ds.drop_vars([\"SPM_QI\", \"TUR_QI\", \"crs\"])\n",
" \n",
"\n",
" def dispose(self):\n",
" self.slice_ds.close()\n",
" self.slice_ds = None\n",
Expand Down Expand Up @@ -5789,31 +5787,34 @@
"source": [
"from zappend.api import SliceSource\n",
"\n",
"\n",
"class MyGrottySliceSource(SliceSource):\n",
" def __init__(self, slice_path):\n",
" self.slice_path = slice_path\n",
" self.slice_ds = None\n",
" self.is_bad = \"20230605\" in slice_path\n",
" self.num_chunks_seen = 0\n",
" \n",
"\n",
" def half_spm(self, spm):\n",
" if self.is_bad:\n",
" self.num_chunks_seen += 1\n",
" if self.num_chunks_seen == 12:\n",
" raise ValueError(\"Bad chunk detected!\")\n",
" return 0.5 * spm\n",
" \n",
" def get_dataset(self) -> xr.Dataset: \n",
"\n",
" def get_dataset(self) -> xr.Dataset:\n",
" LOG.info(f\"Processing slice {self.slice_path}\")\n",
" slice_ds = xr.open_dataset(self.slice_path) \n",
" slice_ds = xr.open_dataset(self.slice_path)\n",
" slice_ds = slice_ds.drop_vars([\"SPM_QI\", \"TUR_QI\", \"crs\"])\n",
" slice_ds = slice_ds.chunk(dict(lon=2000, lat=1000))\n",
" slice_ds = slice_ds.chunk({\"lon\": 2000, \"lat\": 1000})\n",
" for v in slice_ds.data_vars.values():\n",
" v.encoding = {} \n",
" slice_ds[\"SPM_05\"] = slice_ds[\"SPM\"].map_blocks(self.half_spm, template=slice_ds[\"SPM\"]) \n",
" v.encoding = {}\n",
" slice_ds[\"SPM_05\"] = slice_ds[\"SPM\"].map_blocks(\n",
" self.half_spm, template=slice_ds[\"SPM\"]\n",
" )\n",
" self.slice_ds = slice_ds\n",
" return slice_ds\n",
" \n",
"\n",
" def dispose(self):\n",
" self.slice_ds.close()\n",
" self.slice_ds = None\n",
Expand Down
Loading
Loading