Skip to content

Commit ddec00c

Browse files
author
arusu
committed
Containerise echodata and fix the pipeline for non-Azure backends
Everything here was required to run the pipeline outside Azure, on Kubernetes, for the Denoise Lab deployment. Packaging - deploy/Dockerfile.echodata builds the oceanstream-echodata base image (echopype from the oceanstream-integration branch, so the builder needs git for setuptools_scm). boto3 is installed in the same pip run as the echodata extra: resolving it separately lets pip move botocore out from under s3fs, which then fails inside a Dask worker at flow-run time. - .dockerignore keeps data/, raw_data/ and outputs out of the build context; ACR Tasks uploads the whole thing before it starts. Dask worker plugins - S3StoragePlugin and LocalStoragePlugin now subclass WorkerPlugin. distributed 2025.x rejects a plain class with a setup method: "TypeError: Registering duck-typed plugins is not allowed". This broke every non-Azure backend, local included. Memory sizing in a container - effective_parallel_workers sized itself from os.sysconf(SC_PHYS_PAGES), which reports the HOST's RAM inside a container. A pod with a 24 GiB limit on a 110 GiB node picked 7 stage workers and was OOM-killed. Now takes the minimum of that and the cgroup v1/v2 memory limit. - keep_raw was swallowed into the previous line by a literal \n, so it was a comment rather than a dataclass field. Rendering - combine_38khz_day resampled with np.interp, which bridges interior NaN runs. On denoised input those runs ARE the removed noise, so the combined panel painted it back in and looked undenoised. The validity mask is now interpolated alongside and the output blanked where the nearest source samples were mostly NaN. - Combined panels use the same canvas as the per-category ones (34x12 in at 180 dpi) so the two can be compared side by side. - --force no longer reuses existing echogram PNGs. It means the denoise parameters changed, so the previous run's images were being shown for the new parameters, which is silently wrong.
1 parent cafbdd5 commit ddec00c

7 files changed

Lines changed: 270 additions & 16 deletions

File tree

‎.dockerignore‎

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
# =============================================================================
2+
# Build context filter for `az acr build` / `docker build`
3+
# =============================================================================
4+
# ACR Tasks uploads the whole context before it starts, so anything listed here
5+
# is bandwidth saved on every build. `raw_data/` and `data/` in particular sit
6+
# at the repo root on a developer machine and can run to hundreds of GB.
7+
# =============================================================================
8+
9+
# Local data and pipeline outputs — never belongs in an image
10+
data/
11+
raw_data/
12+
out/
13+
**/out/
14+
output/
15+
**/geoparquet/
16+
.oceanstream_staging/
17+
.oceanstream_work/
18+
scripts/batch_processing/raw_cache/
19+
scripts/batch_processing/report_assets/
20+
*.out
21+
22+
# Test fixtures. The image runs the pipeline, it does not run the test suite.
23+
oceanstream/tests/
24+
25+
# Host artifacts that must never overlay what the image builds itself
26+
**/.venv
27+
venv/
28+
venv311/
29+
env/
30+
**/__pycache__
31+
*.py[cod]
32+
*.egg-info/
33+
build/
34+
dist/
35+
.eggs/
36+
37+
# Caches and coverage
38+
.pytest_cache/
39+
.cache/
40+
.coverage
41+
.coverage.*
42+
htmlcov/
43+
coverage_html/
44+
coverage.json
45+
.mypy_cache/
46+
.ruff_cache/
47+
48+
# Docs live in their own repo; notebooks and examples are not runtime code
49+
docs/
50+
docs-internal/
51+
notebooks/
52+
examples/
53+
.docs-website/
54+
.mkdocs-shadcn-fork/
55+
_oceanstream-website-static/
56+
_echodata-legacy-code/
57+
58+
# Secrets and VCS
59+
.env
60+
.env.*
61+
.git/
62+
.github/
63+
64+
# OS / editor cruft
65+
.DS_Store
66+
.AppleDouble
67+
.LSOverride
68+
Thumbs.db
69+
.vscode/
70+
.idea/
71+
.history/
72+
*.bak*
73+
*.log

‎deploy/Dockerfile.echodata‎

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
# =============================================================================
2+
# oceanstream-echodata — base image for the batch echodata pipeline
3+
# =============================================================================
4+
# Ships `oceanstream[echodata]` plus every third-party module the scripts under
5+
# scripts/batch_processing/ import, and the scripts themselves. It runs nothing
6+
# by itself; downstream images (denoise-lab-worker) add their own entrypoint.
7+
#
8+
# Build from the repo root:
9+
# az acr build --registry oceanstreamdevacr --platform linux/amd64 \
10+
# --image oceanstream-echodata:<tag> -f deploy/Dockerfile.echodata .
11+
# =============================================================================
12+
13+
ARG PYTHON_VERSION=3.12
14+
15+
# -----------------------------------------------------------------------------
16+
# Stage 1: build the virtualenv
17+
# -----------------------------------------------------------------------------
18+
# Pinned to bookworm: trixie renamed the GDAL runtime lib (libgdal32 ->
19+
# libgdal36), so the runtime stage below could not satisfy it.
20+
FROM python:${PYTHON_VERSION}-slim-bookworm AS builder
21+
22+
# git is not optional: pyproject.toml pulls echopype from the OceanStreamIO
23+
# fork's `oceanstream-integration` branch, and it builds with setuptools_scm,
24+
# which reads the cloned .git to derive a version.
25+
RUN apt-get update && apt-get install -y --no-install-recommends \
26+
build-essential \
27+
git \
28+
libgdal-dev \
29+
libgeos-dev \
30+
libproj-dev \
31+
&& rm -rf /var/lib/apt/lists/*
32+
33+
RUN python -m venv /opt/venv
34+
ENV PATH="/opt/venv/bin:$PATH"
35+
RUN pip install --no-cache-dir --upgrade pip wheel
36+
37+
WORKDIR /src
38+
39+
# Resolve the dependency tree against a package skeleton first. Editing the
40+
# source then only re-runs the --no-deps install below, instead of re-resolving
41+
# echopype-from-git, dask[complete], copernicusmarine and friends.
42+
COPY pyproject.toml ./
43+
COPY oceanstream/README.md ./oceanstream/README.md
44+
45+
# The trailing names are modules the batch scripts import that the `echodata`
46+
# extra does not cover. They go in the SAME pip invocation on purpose: s3fs
47+
# pulls aiobotocore, which pins botocore to a narrow range, so boto3 has to be
48+
# resolved alongside it rather than bolted on afterwards.
49+
# Deliberately absent is the `echodata-viz` extra — holoviews/hvplot/panel/
50+
# bokeh/datashader are reached only by plot_interactive_echogram(), which
51+
# imports them lazily and which this pipeline never calls.
52+
RUN touch oceanstream/__init__.py \
53+
&& pip install --no-cache-dir ".[echodata]" \
54+
boto3 \
55+
cartopy \
56+
cmocean \
57+
matplotlib \
58+
rasterio \
59+
geopandas \
60+
scipy \
61+
pillow \
62+
psycopg2-binary \
63+
azure-storage-file-share \
64+
python-dotenv
65+
66+
COPY oceanstream/ ./oceanstream/
67+
RUN pip install --no-cache-dir --no-deps --force-reinstall .
68+
69+
# -----------------------------------------------------------------------------
70+
# Stage 2: runtime
71+
# -----------------------------------------------------------------------------
72+
# Must match the builder's Debian release — the venv links against its libs.
73+
FROM python:${PYTHON_VERSION}-slim-bookworm AS runtime
74+
75+
LABEL org.opencontainers.image.title="oceanstream-echodata"
76+
LABEL org.opencontainers.image.description="oceanstream[echodata] plus the batch_processing pipeline scripts"
77+
LABEL org.opencontainers.image.source="https://github.com/OceanStreamIO/oceanstream-cli"
78+
79+
RUN apt-get update && apt-get install -y --no-install-recommends \
80+
libgdal32 \
81+
libgeos-c1v5 \
82+
libproj25 \
83+
ca-certificates \
84+
curl \
85+
&& rm -rf /var/lib/apt/lists/* \
86+
&& apt-get clean
87+
88+
# uid/gid 1000 is fixed on purpose: the worker and web containers share PVCs,
89+
# and the pod's fsGroup has to match one number across every image.
90+
RUN groupadd -g 1000 oceanstream \
91+
&& useradd -u 1000 -g 1000 -m -d /home/oceanstream -s /usr/sbin/nologin oceanstream
92+
93+
COPY --from=builder /opt/venv /opt/venv
94+
COPY scripts/ /opt/oceanstream/scripts/
95+
96+
# MPLBACKEND: without Agg, matplotlib picks a GUI backend and dies on import.
97+
# MPLCONFIGDIR/XDG_CACHE_HOME default under $HOME, which is not writable once
98+
# the pod runs as an arbitrary fsGroup; /tmp always is.
99+
ENV PATH="/opt/venv/bin:$PATH" \
100+
PYTHONUNBUFFERED=1 \
101+
PYTHONDONTWRITEBYTECODE=1 \
102+
HOME=/home/oceanstream \
103+
OCEANSTREAM_SCRIPTS=/opt/oceanstream/scripts/batch_processing \
104+
MPLBACKEND=Agg \
105+
MPLCONFIGDIR=/tmp/matplotlib \
106+
XDG_CACHE_HOME=/tmp/cache
107+
108+
WORKDIR /opt/oceanstream/scripts/batch_processing
109+
USER oceanstream
110+
111+
CMD ["python"]

‎oceanstream/echodata/plot/combined.py‎

Lines changed: 25 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -55,8 +55,13 @@
5555
# Sv panel look indistinguishable from a gridded MVBS one.
5656
PINGS_PER_PIXEL: float = 2.0
5757

58-
# Cap on figure width (inches). At dpi=150 this is 15 000 px.
59-
MAX_WIDTH_IN: float = 100.0
58+
# Canvas geometry, matched to the per-category echograms in echogram.py
59+
# (34 x 12 in at 180 dpi = 6120 x 2160 px) so the two can be compared side by
60+
# side. Honouring PINGS_PER_PIXEL for a 29 000-ping day would want 80 in, and
61+
# the resulting 14 500 x 1050 px strip is unreadable in a browser.
62+
MAX_WIDTH_IN: float = 34.0
63+
HEIGHT_IN: float = 12.0
64+
RENDER_DPI: int = 180
6065

6166
# Pulse-mode indicator bar colours (matches legacy scripts)
6267
_PULSE_COLORS = {
@@ -192,10 +197,20 @@ def combine_38khz_day(
192197
row = sv_raw[i, valid]
193198
m = ~np.isnan(row)
194199
if m.sum() > 1:
195-
sv_interp[i] = np.interp(
200+
interp = np.interp(
196201
common_depth, depth_valid[m], row[m],
197202
left=np.nan, right=np.nan,
198203
)
204+
# np.interp bridges interior NaN runs, and on a denoised input
205+
# those runs ARE the removed noise — bridging them paints it
206+
# back in and the panel looks undenoised. Blank any output
207+
# sample whose nearest source samples were mostly NaN.
208+
kept = np.interp(
209+
common_depth, depth_valid, m.astype(np.float32),
210+
left=0.0, right=0.0,
211+
)
212+
interp[kept < 0.5] = np.nan # noqa: PLR2004 - nearest-sample majority
213+
sv_interp[i] = interp
199214

200215
mode_code = 0 if mode == "long_pulse" else 1
201216
ds_new = xr.Dataset(
@@ -321,9 +336,10 @@ def render_combined_echogram(
321336
max_plot_depth: float = MAX_PLOT_DEPTH,
322337
vmin: float = SV_VMIN,
323338
vmax: float = SV_VMAX,
324-
dpi: int = 150,
339+
dpi: int = RENDER_DPI,
325340
pings_per_pixel: float = PINGS_PER_PIXEL,
326341
max_width_in: float = MAX_WIDTH_IN,
342+
height_in: float = HEIGHT_IN,
327343
) -> Optional[Path]:
328344
"""Render a combined 24h echogram for the 38 kHz channel.
329345
@@ -355,11 +371,12 @@ def render_combined_echogram(
355371
max_plot_depth : float
356372
Trim depth axis at this value (metres).
357373
pings_per_pixel : float
358-
Target pings per horizontal pixel. The canvas is widened until this
359-
density is met, so a 29 000-ping Sv panel gets a much wider figure
360-
than an 8 500-bin MVBS one.
374+
Target pings per horizontal pixel. The canvas is widened towards this
375+
density until ``max_width_in`` stops it.
361376
max_width_in : float
362377
Upper bound on figure width (inches).
378+
height_in : float
379+
Figure height (inches), including the pulse-mode strip.
363380
364381
Returns
365382
-------
@@ -418,7 +435,7 @@ def render_combined_echogram(
418435
max(14.0, time_span * 1.2, n_pings / max(pings_per_pixel * dpi, 1.0)),
419436
)
420437
)
421-
height = 7.0
438+
height = height_in
422439

423440
has_pulse = pulse_mode is not None
424441
cbar_frac = 0.3 / width

‎scripts/batch_processing/config.py‎

Lines changed: 34 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,29 @@
1313
from typing import Optional
1414

1515

16+
def _cgroup_memory_limit_gb() -> Optional[float]:
17+
"""The container's memory ceiling, or None when unconfined."""
18+
for path in (
19+
"/sys/fs/cgroup/memory.max", # cgroup v2
20+
"/sys/fs/cgroup/memory/memory.limit_in_bytes", # cgroup v1
21+
):
22+
try:
23+
raw = Path(path).read_text().strip()
24+
except OSError:
25+
continue
26+
if raw == "max":
27+
return None
28+
try:
29+
value = int(raw)
30+
except ValueError:
31+
continue
32+
# cgroup v1 reports a sentinel near 2**63 when there is no limit.
33+
if value >= 1 << 62:
34+
return None
35+
return value / (1024 ** 3)
36+
return None
37+
38+
1639
@dataclass
1740
class DaskConfig:
1841
"""Dask distributed cluster settings."""
@@ -323,7 +346,8 @@ class PipelineConfig:
323346
skip_campaign_echograms: bool = False # skip only the campaign echogram loop (keep campaign zarr)
324347
build_campaign_sv_zarr: bool = False # experimental
325348
category_parallel: bool = True # parallelize short_pulse/long_pulse within each day
326-
resume_stage: int = 0 # resume from this stage (0 = start from beginning)\n keep_raw: bool = False # keep downloaded raw files after conversion
349+
resume_stage: int = 0 # resume from this stage (0 = start from beginning)
350+
keep_raw: bool = False # keep downloaded raw files after conversion
327351
# Stop after this stage (0 = run everything). Single-day comparison runs
328352
# set this to 9 so the campaign-aggregation stages don't run on one day.
329353
stop_after_stage: int = 0
@@ -395,16 +419,22 @@ def from_env(cls) -> PipelineConfig:
395419
def effective_parallel_workers(self, mem_per_worker_gb: float = 2.0) -> int:
396420
"""Return the number of parallel stage workers.
397421
398-
If ``parallel_workers`` is 0 (the default), auto-detect from total
399-
system RAM, reserving memory for Dask workers and OS overhead.
400-
Each parallel denoise/echogram task needs ~2 GB.
422+
If ``parallel_workers`` is 0 (the default), auto-detect from the RAM
423+
this process may actually use, reserving memory for Dask workers and
424+
OS overhead. Each parallel denoise/echogram task needs ~2 GB.
401425
"""
402426
if self.parallel_workers > 0:
403427
return self.parallel_workers
404428
try:
405429
total_gb = os.sysconf("SC_PAGE_SIZE") * os.sysconf("SC_PHYS_PAGES") / (1024 ** 3)
406430
except (ValueError, OSError):
407431
return 2 # safe default if detection fails
432+
# sysconf reports the HOST's RAM even inside a container, so a pod with
433+
# a 24 GiB limit on a 256 GiB node would size itself for 256 GiB and be
434+
# OOM-killed or evicted mid-run.
435+
cgroup_gb = _cgroup_memory_limit_gb()
436+
if cgroup_gb is not None:
437+
total_gb = min(total_gb, cgroup_gb)
408438
# Reserve ~50% for the Dask LocalCluster, OS, and headroom
409439
available_gb = total_gb * 0.5
410440
workers = max(1, int(available_gb / mem_per_worker_gb))

‎scripts/batch_processing/local_storage.py‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -216,11 +216,21 @@ def ls(self, path: str, detail: bool = False) -> list:
216216

217217
# ── Dask worker plugin ───────────────────────────────────────────────────
218218

219-
class LocalStoragePlugin:
219+
try: # distributed is only needed for the Dask path, not for patch_storage()
220+
from distributed.diagnostics.plugin import WorkerPlugin as _WorkerPlugin
221+
except ImportError:
222+
_WorkerPlugin = object
223+
224+
225+
class LocalStoragePlugin(_WorkerPlugin):
220226
"""Dask worker plugin that applies local storage patches on each worker.
221227
222228
Register with ``client.register_worker_plugin(LocalStoragePlugin(root))``
223229
so that Dask workers use local filesystem instead of Azure.
230+
231+
Subclassing ``WorkerPlugin`` is mandatory, not decorative: distributed
232+
2025.x onwards raises ``TypeError: Registering duck-typed plugins is not
233+
allowed`` for a plain class with a ``setup`` method.
224234
"""
225235

226236
name = "local-storage"

‎scripts/batch_processing/process_campaign.py‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2237,8 +2237,11 @@ def run_echogram_generation(
22372237
with ThreadPoolExecutor(max_workers=max_workers) as pool:
22382238
for day_key in sorted(all_keys):
22392239
for category in sorted(all_keys[day_key]):
2240-
# Skip if echograms already exist in Azure
2241-
existing = _count_existing_echograms(output_container, day_key, category)
2240+
# --force means the denoise parameters changed, so the existing
2241+
# PNGs render the previous run and must not be reused.
2242+
existing = (
2243+
0 if cfg.force else _count_existing_echograms(output_container, day_key, category)
2244+
)
22422245
if existing > 0:
22432246
logger.info(" Skipping %s/%s — %d echograms already exist", day_key, category, existing)
22442247
continue

‎scripts/batch_processing/s3_storage.py‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -363,13 +363,23 @@ def get_mapper(self, path: str, **kwargs) -> fsspec.FSMap:
363363

364364
# ── Dask worker plugin ───────────────────────────────────────────────────
365365

366-
class S3StoragePlugin:
366+
try: # distributed is only needed for the Dask path, not for patch_storage()
367+
from distributed.diagnostics.plugin import WorkerPlugin as _WorkerPlugin
368+
except ImportError:
369+
_WorkerPlugin = object
370+
371+
372+
class S3StoragePlugin(_WorkerPlugin):
367373
"""Dask worker plugin that applies the S3 storage patches on each worker.
368374
369375
Register with ``client.register_plugin(S3StoragePlugin(bucket, prefix))``.
370376
Credentials are deliberately not carried by default: leaving *key* and
371377
*secret* unset makes each worker read them from its own environment
372378
(a K8s secret) instead of shipping them through the scheduler.
379+
380+
Subclassing ``WorkerPlugin`` is mandatory, not decorative: distributed
381+
2025.x onwards raises ``TypeError: Registering duck-typed plugins is not
382+
allowed`` for a plain class with a ``setup`` method.
373383
"""
374384

375385
name = "s3-storage"

0 commit comments

Comments
 (0)