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
45 changes: 45 additions & 0 deletions docs/Compatibility_Fixes_2026-02-13.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
# Compatibility Fixes (2026-02-13)

This note summarizes compatibility hardening completed for recent pandas/numpy/joblib behavior changes.

## Core runtime fixes

- Added safe parallel worker resolution helper: `safe_n_jobs(task_count, thread_count)` in `pingmapper/funcs_common.py`.
- Replaced fragile `Parallel(n_jobs=np.min([len(...), threadCnt]))` patterns in core runtime modules.
- Fixed pandas chained assignment hotspots to direct `.loc[row_mask, col] = value` writes.
- Replaced positional-usage `loc` with `iloc` where label-based indexing could fail.
- Restored robust heading filtering and chunk reassignment logic in `pingmapper/class_sonObj.py`:
- Heading filter now handles full-window spread and end-window coverage.
- Empty-data guards prevent `IndexError` in chunk reassignment.

## Utility/workflow fixes

- Applied safe `n_jobs` handling to utility workflows and draft scripts.
- Updated positional `.loc[0]` usages to `.iloc[0]` in summary workflows where row-position semantics were intended.

## Key files updated

- `pingmapper/funcs_common.py`
- `pingmapper/main_readFiles.py`
- `pingmapper/main_mapSubstrate.py`
- `pingmapper/main_rectify.py`
- `pingmapper/funcs_rectify.py`
- `pingmapper/class_portstarObj.py`
- `pingmapper/class_sonObj.py`
- `pingmapper/class_sonObj_nadirgaptest.py`
- `pingmapper/utils/main_mosaic_transects.py`
- `pingmapper/utils/RawEGN_avg_predictions.py`
- `pingmapper/utils/DRAFT_Workflows/avg_predictions_Mussel_WBL.py`
- `pingmapper/utils/Substrate_Summaries/summarize_project_substrate.py`
- `pingmapper/utils/Substrate_Summaries/02_gen_summary_stamp_shps.py`

## Validation performed

- Read-stage smoke test passed on:
- `Z:\miniforge3\envs\ping\exampleData\Test-Small-DS.DAT`
- Full end-to-end smoke test passed (read + rectify + map) on the same dataset.

## Environment note

- In this shell context, `conda run -p Z:\miniforge3\envs\ping --no-capture-output python ...` was reliable.
- Direct `Z:\miniforge3\envs\ping\python.exe` invocation showed interpreter crash behavior in this session context.
43 changes: 23 additions & 20 deletions pingmapper/class_portstarObj.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,10 +59,13 @@

import inspect

quiet_tensorflow_warnings()

try:
from doodleverse_utils.imports import *
from doodleverse_utils.model_imports import *
from doodleverse_utils.prediction_imports import *
with suppress_stdout_stderr():
from doodleverse_utils.imports import *
from doodleverse_utils.model_imports import *
from doodleverse_utils.prediction_imports import *
except ImportError as e:
import traceback
print('\n' + '='*80)
Expand Down Expand Up @@ -348,35 +351,35 @@ def _createMosaic(self,
if mosaic == 1:
if son:
if self.port.rect_wcp:
_ = Parallel(n_jobs= np.min([len(wcpToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(wcpToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
if self.port.rect_wcr:
_ = Parallel(n_jobs= np.min([len(srcToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(srcToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
else:
if self.port.map_sub:
_ = Parallel(n_jobs= np.min([len(subToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([sub], overview=overview, i=i, son=son) for i, sub in enumerate(subToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(subToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([sub], overview=overview, i=i, son=son) for i, sub in enumerate(subToMosaic))

if self.port.map_predict:
# Determine number of bands, i.e. substrate classes
bands = self._getBandCount(predictToMosaic[0][0])
for i, pred in enumerate(predictToMosaic):
_ = Parallel(n_jobs= np.min([bands, threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))
_ = Parallel(n_jobs=safe_n_jobs(bands, threadCnt), verbose=10)(delayed(self._mosaicGtiff)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))

# Create vrt
elif mosaic == 2:
if son:
if self.port.rect_wcp:
_ = Parallel(n_jobs= np.min([len(wcpToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(wcpToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
if self.port.rect_wcr:
_ = Parallel(n_jobs= np.min([len(srcToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(srcToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
else:
if self.port.map_sub:
_ = Parallel(n_jobs= np.min([len(subToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([sub], overview, i, son=son) for i, sub in enumerate(subToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(subToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([sub], overview, i, son=son) for i, sub in enumerate(subToMosaic))

if self.port.map_predict:
# Determine number of bands, i.e. substrate classes
bands = self._getBandCount(predictToMosaic[0][0])
for i, pred in enumerate(predictToMosaic):
_ = Parallel(n_jobs= np.min([bands, threadCnt]), verbose=10)(delayed(self._mosaicVRT)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))
_ = Parallel(n_jobs=safe_n_jobs(bands, threadCnt), verbose=10)(delayed(self._mosaicVRT)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))

return

Expand Down Expand Up @@ -549,35 +552,35 @@ def _createMosaicTransect(self,
if mosaic == 1:
if son:
if self.port.rect_wcp:
_ = Parallel(n_jobs= np.min([len(wcpToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(wcpToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
if self.port.rect_wcr:
_ = Parallel(n_jobs= np.min([len(srcToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(srcToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
else:
if self.port.map_sub:
_ = Parallel(n_jobs= np.min([len(subToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([sub], overview=overview, i=i, son=son) for i, sub in enumerate(subToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(subToMosaic), threadCnt), verbose=10)(delayed(self._mosaicGtiff)([sub], overview=overview, i=i, son=son) for i, sub in enumerate(subToMosaic))

if self.port.map_predict:
# Determine number of bands, i.e. substrate classes
bands = self._getBandCount(predictToMosaic[0][0])
for i, pred in enumerate(predictToMosaic):
_ = Parallel(n_jobs= np.min([bands, threadCnt]), verbose=10)(delayed(self._mosaicGtiff)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))
_ = Parallel(n_jobs=safe_n_jobs(bands, threadCnt), verbose=10)(delayed(self._mosaicGtiff)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))

# Create vrt
elif mosaic == 2:
if son:
if self.port.rect_wcp:
_ = Parallel(n_jobs= np.min([len(wcpToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(wcpToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([wcp], overview, i, son=son) for i, wcp in enumerate(wcpToMosaic))
if self.port.rect_wcr:
_ = Parallel(n_jobs= np.min([len(srcToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(srcToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([src], overview, i, son=son) for i, src in enumerate(srcToMosaic))
else:
if self.port.map_sub:
_ = Parallel(n_jobs= np.min([len(subToMosaic), threadCnt]), verbose=10)(delayed(self._mosaicVRT)([sub], overview, i, son=son) for i, sub in enumerate(subToMosaic))
_ = Parallel(n_jobs=safe_n_jobs(len(subToMosaic), threadCnt), verbose=10)(delayed(self._mosaicVRT)([sub], overview, i, son=son) for i, sub in enumerate(subToMosaic))

if self.port.map_predict:
# Determine number of bands, i.e. substrate classes
bands = self._getBandCount(predictToMosaic[0][0])
for i, pred in enumerate(predictToMosaic):
_ = Parallel(n_jobs= np.min([bands, threadCnt]), verbose=10)(delayed(self._mosaicVRT)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))
_ = Parallel(n_jobs=safe_n_jobs(bands, threadCnt), verbose=10)(delayed(self._mosaicVRT)([pred], overview, i, bands=[c], son=True) for c in range(1,bands+1))

return

Expand Down Expand Up @@ -2572,7 +2575,7 @@ def _rasterToPoly(self, mosaic, threadCnt, mosaic_nchunk):
os.mkdir(outDir)

print("\n\tExporting to shapefile...")
_ = Parallel(n_jobs= np.min([len(rasterFiles), threadCnt]), verbose=10)(delayed(self._createPolygon)(f, outDir) for f in rasterFiles)
_ = Parallel(n_jobs=safe_n_jobs(len(rasterFiles), threadCnt), verbose=10)(delayed(self._createPolygon)(f, outDir) for f in rasterFiles)

return

Expand Down
71 changes: 55 additions & 16 deletions pingmapper/class_sonObj.py
Original file line number Diff line number Diff line change
Expand Up @@ -332,32 +332,63 @@ def _filterHeading(self,
dist_start = 0 # Counter for beginning of current window
dist_end = dist_start + d # Counbter for end of current window

df = df.copy()
df[filtCol] = False

# If distance is invalid, do not remove all pings
if not np.isfinite(max_dist):
df[filtCol] = True
return df

# If heading distance window covers the entire recording,
# evaluate the complete record once.
if max_dist <= d:
head_vals = df[head].to_numpy(copy=True)
head_vals = head_vals[np.isfinite(head_vals)]

if len(head_vals) < 2:
df[filtCol] = True
return df

head_vals = np.deg2rad(head_vals)
head_vals = np.unwrap(head_vals)
vessel_dev = np.ptp(head_vals)

if vessel_dev < dev:
df[filtCol] = True

return df

##############################
# Iterator through each window

# Compare heading deviation from first and last ping for current window
while dist_end < max_dist:
# Compare heading deviation across each window
while dist_start < max_dist:

# Filter df by window
dfFilt = df[(df[trk_dist] >= dist_start) & (df[trk_dist] < dist_end)]
if dist_end < max_dist:
window_mask = (df[trk_dist] >= dist_start) & (df[trk_dist] < dist_end)
else:
window_mask = (df[trk_dist] >= dist_start) & (df[trk_dist] <= max_dist)
dfFilt = df.loc[window_mask]

dfFilt[head] = np.deg2rad(dfFilt[head])
# Unwrap the heading because heading is circular
dfFilt[head] = np.unwrap(dfFilt[head])
if len(dfFilt) > 1:

if len(dfFilt) > 0:
head_vals = dfFilt[head].to_numpy(copy=True)
head_vals = head_vals[np.isfinite(head_vals)]

# Get difference between start and end heading
start = dfFilt[head].iloc[0]
end = dfFilt[head].iloc[-1]
vessel_dev = np.abs(start - end)
if len(head_vals) > 1:
head_vals = np.deg2rad(head_vals)
# Unwrap the heading because heading is circular
head_vals = np.unwrap(head_vals)

# Compare vessel deviation to threshold deviation
if vessel_dev < dev:
# Keep these pings
df[filtCol].loc[dfFilt.index] = True
# Get total heading spread in the window
vessel_dev = np.ptp(head_vals)

# Compare vessel deviation to threshold deviation
if vessel_dev < dev:
# Keep these pings
df.loc[dfFilt.index, filtCol] = True

# dist_start += win
# dist_start = dist_end
Expand Down Expand Up @@ -528,6 +559,12 @@ def _reassignChunks(self,
# Reassign Chunks
nchunk = self.nchunk

if sonDF.empty:
sonDF = sonDF.copy()
sonDF['transect'] = pd.Series(dtype='int64')
sonDF['chunk_id'] = pd.Series(dtype='int64')
return sonDF

# Make transects from consective pings using dataframe index
idx = sonDF.index.values
transect_groups = np.split(idx, np.where(np.diff(idx) != 1)[0]+1)
Expand All @@ -536,7 +573,9 @@ def _reassignChunks(self,
# Assign transect
transect = 0
for t in transect_groups:
sonDF.loc[sonDF.index>=t[0], 'transect'] = transect
if len(t) == 0:
continue
sonDF.loc[t, 'transect'] = transect
transect += 1

# Set chunks
Expand Down
71 changes: 55 additions & 16 deletions pingmapper/class_sonObj_nadirgaptest.py
Original file line number Diff line number Diff line change
Expand Up @@ -332,32 +332,63 @@ def _filterHeading(self,
dist_start = 0 # Counter for beginning of current window
dist_end = dist_start + d # Counbter for end of current window

df = df.copy()
df[filtCol] = False

# If distance is invalid, do not remove all pings
if not np.isfinite(max_dist):
df[filtCol] = True
return df

# If heading distance window covers the entire recording,
# evaluate the complete record once.
if max_dist <= d:
head_vals = df[head].to_numpy(copy=True)
head_vals = head_vals[np.isfinite(head_vals)]

if len(head_vals) < 2:
df[filtCol] = True
return df

head_vals = np.deg2rad(head_vals)
head_vals = np.unwrap(head_vals)
vessel_dev = np.ptp(head_vals)

if vessel_dev < dev:
df[filtCol] = True

return df

##############################
# Iterator through each window

# Compare heading deviation from first and last ping for current window
while dist_end < max_dist:
# Compare heading deviation across each window
while dist_start < max_dist:

# Filter df by window
dfFilt = df[(df[trk_dist] >= dist_start) & (df[trk_dist] < dist_end)]
if dist_end < max_dist:
window_mask = (df[trk_dist] >= dist_start) & (df[trk_dist] < dist_end)
else:
window_mask = (df[trk_dist] >= dist_start) & (df[trk_dist] <= max_dist)
dfFilt = df.loc[window_mask]

dfFilt[head] = np.deg2rad(dfFilt[head])
# Unwrap the heading because heading is circular
dfFilt[head] = np.unwrap(dfFilt[head])
if len(dfFilt) > 1:

if len(dfFilt) > 0:
head_vals = dfFilt[head].to_numpy(copy=True)
head_vals = head_vals[np.isfinite(head_vals)]

# Get difference between start and end heading
start = dfFilt[head].iloc[0]
end = dfFilt[head].iloc[-1]
vessel_dev = np.abs(start - end)
if len(head_vals) > 1:
head_vals = np.deg2rad(head_vals)
# Unwrap the heading because heading is circular
head_vals = np.unwrap(head_vals)

# Compare vessel deviation to threshold deviation
if vessel_dev < dev:
# Keep these pings
df[filtCol].loc[dfFilt.index] = True
# Get total heading spread in the window
vessel_dev = np.ptp(head_vals)

# Compare vessel deviation to threshold deviation
if vessel_dev < dev:
# Keep these pings
df.loc[dfFilt.index, filtCol] = True

# dist_start += win
# dist_start = dist_end
Expand Down Expand Up @@ -528,6 +559,12 @@ def _reassignChunks(self,
# Reassign Chunks
nchunk = self.nchunk

if sonDF.empty:
sonDF = sonDF.copy()
sonDF['transect'] = pd.Series(dtype='int64')
sonDF['chunk_id'] = pd.Series(dtype='int64')
return sonDF

# Make transects from consective pings using dataframe index
idx = sonDF.index.values
transect_groups = np.split(idx, np.where(np.diff(idx) != 1)[0]+1)
Expand All @@ -536,7 +573,9 @@ def _reassignChunks(self,
# Assign transect
transect = 0
for t in transect_groups:
sonDF.loc[sonDF.index>=t[0], 'transect'] = transect
if len(t) == 0:
continue
sonDF.loc[t, 'transect'] = transect
transect += 1

# Set chunks
Expand Down
Loading
Loading