diff --git a/docs/Compatibility_Fixes_2026-02-13.md b/docs/Compatibility_Fixes_2026-02-13.md new file mode 100644 index 0000000..a0c4e14 --- /dev/null +++ b/docs/Compatibility_Fixes_2026-02-13.md @@ -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. \ No newline at end of file diff --git a/pingmapper/class_portstarObj.py b/pingmapper/class_portstarObj.py index 0769c79..c85398d 100644 --- a/pingmapper/class_portstarObj.py +++ b/pingmapper/class_portstarObj.py @@ -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) @@ -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 @@ -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 @@ -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 diff --git a/pingmapper/class_sonObj.py b/pingmapper/class_sonObj.py index fdc3d4f..713e681 100644 --- a/pingmapper/class_sonObj.py +++ b/pingmapper/class_sonObj.py @@ -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 @@ -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) @@ -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 diff --git a/pingmapper/class_sonObj_nadirgaptest.py b/pingmapper/class_sonObj_nadirgaptest.py index 54e8b41..31b2ade 100644 --- a/pingmapper/class_sonObj_nadirgaptest.py +++ b/pingmapper/class_sonObj_nadirgaptest.py @@ -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 @@ -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) @@ -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 diff --git a/pingmapper/funcs_common.py b/pingmapper/funcs_common.py index abd2925..5f4e256 100644 --- a/pingmapper/funcs_common.py +++ b/pingmapper/funcs_common.py @@ -29,7 +29,7 @@ # OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE # SOFTWARE. -import os, sys, struct, gc +import os, sys, struct, gc, io, contextlib # Add 'pingmapper' to the path, may not need after pypi package... SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__)) @@ -96,6 +96,78 @@ # from funcs_pyhum_correct import doPyhumCorrections +# ========================================================= +def safe_n_jobs(task_count, thread_count=0): + ''' + Resolve a valid joblib n_jobs value from task count and user thread setting. + Always returns at least 1 to avoid joblib ValueError on n_jobs == 0. + ''' + try: + task_count = int(task_count) + except Exception: + task_count = 0 + + try: + thread_count = int(thread_count) + except Exception: + thread_count = 0 + + if thread_count <= 0: + thread_count = cpu_count() + + if thread_count < 1: + thread_count = 1 + + if task_count < 1: + return 1 + + return int(min(task_count, thread_count)) + + +# ========================================================= +def quiet_tensorflow_warnings(): + ''' + Reduce TensorFlow/absl informational and warning output. + Safe to call even if TensorFlow is not installed. + ''' + os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' + os.environ['AUTOGRAPH_VERBOSITY'] = '0' + + logging.getLogger('tensorflow').setLevel(logging.ERROR) + logging.getLogger('absl').setLevel(logging.ERROR) + + try: + import absl.logging as absl_logging + absl_logging.set_verbosity(absl_logging.ERROR) + absl_logging.set_stderrthreshold('error') + except Exception: + pass + + try: + import tensorflow as tf + tf.get_logger().setLevel('ERROR') + try: + tf.compat.v1.logging.set_verbosity(tf.compat.v1.logging.ERROR) + except Exception: + pass + try: + tf.autograph.set_verbosity(0) + except Exception: + pass + except Exception: + pass + + +# ========================================================= +@contextlib.contextmanager +def suppress_stdout_stderr(): + ''' + Context manager to suppress stdout/stderr noise from third-party imports. + ''' + with contextlib.redirect_stdout(io.StringIO()), contextlib.redirect_stderr(io.StringIO()): + yield + + # ========================================================= def rescale( dat, mn, diff --git a/pingmapper/funcs_model.py b/pingmapper/funcs_model.py index e727d42..294dab5 100644 --- a/pingmapper/funcs_model.py +++ b/pingmapper/funcs_model.py @@ -39,7 +39,7 @@ sys.path.append(PACKAGE_DIR) from pingmapper.funcs_common import * -os.environ['TF_CPP_MIN_LOG_LEVEL']='3' +quiet_tensorflow_warnings() import json import numpy as np # import tensorflow as tf @@ -72,11 +72,12 @@ logging.set_verbosity_error() # Fixes depth detection warning - tf.get_logger().setLevel('ERROR') + quiet_tensorflow_warnings() - 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 * DEPTH_DETECTION_AVAILABLE = True except ImportError as e: diff --git a/pingmapper/funcs_rectify.py b/pingmapper/funcs_rectify.py index 07677c8..5a0190c 100644 --- a/pingmapper/funcs_rectify.py +++ b/pingmapper/funcs_rectify.py @@ -168,8 +168,8 @@ def smoothTrackline(projDir='', x_offset='', y_offset='', nchunk ='', cog=True, for name, group in sonDF.groupby('transect'): for n, g in group.groupby('chunk_id'): # sonDF.loc[sonDF['chunk_id'] == n] - sonDF['chunk_id'].loc[sonDF['chunk_id'] == n] = c - sonDF['transect'].loc[sonDF['transect'] == name] = t + sonDF.loc[sonDF['chunk_id'] == n, 'chunk_id'] = c + sonDF.loc[sonDF['transect'] == name, 'transect'] = t c += 1 t+=1 @@ -212,8 +212,8 @@ def smoothTrackline(projDir='', x_offset='', y_offset='', nchunk ='', cog=True, for name, group in sonDF.groupby('transect'): for n, g in group.groupby('chunk_id'): # sonDF.loc[sonDF['chunk_id'] == n] - sonDF['chunk_id'].loc[sonDF['chunk_id'] == n] = c - sonDF['transect'].loc[sonDF['transect'] == name] = t + sonDF.loc[sonDF['chunk_id'] == n, 'chunk_id'] = c + sonDF.loc[sonDF['transect'] == name, 'transect'] = t c += 1 t+=1 @@ -311,7 +311,7 @@ def smoothTrackline(projDir='', x_offset='', y_offset='', nchunk ='', cog=True, print("\nCalculating, smoothing, and interpolating range extent coordinates...") else: print("\nCalculating range extent coordinates from vessel heading...") - Parallel(n_jobs= np.min([len(portstar), threadCnt]), verbose=10)(delayed(son._getRangeCoords)(flip, filterRange, cog) for son in portstar) + Parallel(n_jobs=safe_n_jobs(len(portstar), threadCnt), verbose=10)(delayed(son._getRangeCoords)(flip, filterRange, cog) for son in portstar) print("Done!") print("Time (s):", round(time.time() - start_time, ndigits=1)) gc.collect() diff --git a/pingmapper/main_mapSubstrate.py b/pingmapper/main_mapSubstrate.py index c683a9a..4dea177 100644 --- a/pingmapper/main_mapSubstrate.py +++ b/pingmapper/main_mapSubstrate.py @@ -272,7 +272,7 @@ def map_master_func(logfilename='', # Do prediction (make parallel later) print('\n\tPredicting substrate for', len(chunks), son.beamName, 'chunks') - Parallel(n_jobs=np.min([len(chunks), threadCnt]))(delayed(son._detectSubstrate)(i, USE_GPU) for i in tqdm(chunks)) + Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._detectSubstrate)(i, USE_GPU) for i in tqdm(chunks)) son._cleanup() son._pickleSon() @@ -321,7 +321,7 @@ def map_master_func(logfilename='', # Plot substrate classification() # sys.exit() - Parallel(n_jobs=np.min([len(toMap), threadCnt]))(delayed(son._pltSubClass)(map_class_method, c, f, spdCor=spdCor, maxCrop=maxCrop, probs=probs) for c, f in tqdm((toMap.items()))) + Parallel(n_jobs=safe_n_jobs(len(toMap), threadCnt))(delayed(son._pltSubClass)(map_class_method, c, f, spdCor=spdCor, maxCrop=maxCrop, probs=probs) for c, f in tqdm((toMap.items()))) son._pickleSon() del toMap @@ -381,7 +381,7 @@ def map_master_func(logfilename='', # Create portstarObj psObj = portstarObj(mapObjs) - Parallel(n_jobs=np.min([len(toMap), threadCnt]))(delayed(psObj._mapSubstrate)(map_class_method, c, f) for c, f in tqdm(toMap.items())) + Parallel(n_jobs=safe_n_jobs(len(toMap), threadCnt))(delayed(psObj._mapSubstrate)(map_class_method, c, f) for c, f in tqdm(toMap.items())) del toMap print("\nDone!") @@ -521,7 +521,7 @@ def map_master_func(logfilename='', # Create portstarObj psObj = portstarObj(mapObjs) - Parallel(n_jobs=np.min([len(toMap), threadCnt]))(delayed(psObj._mapPredictions)(map_predict, 'map_'+a, c, f) for c, f in tqdm(toMap.items())) + Parallel(n_jobs=safe_n_jobs(len(toMap), threadCnt))(delayed(psObj._mapPredictions)(map_predict, 'map_'+a, c, f) for c, f in tqdm(toMap.items())) del toMap, psObj print("\nDone!") diff --git a/pingmapper/main_readFiles.py b/pingmapper/main_readFiles.py index c7ac3e3..b71e94b 100644 --- a/pingmapper/main_readFiles.py +++ b/pingmapper/main_readFiles.py @@ -43,8 +43,11 @@ import shutil +quiet_tensorflow_warnings() + try: - from doodleverse_utils.imports import * + with suppress_stdout_stderr(): + from doodleverse_utils.imports import * except ImportError as e: import traceback print('\n' + '='*80) @@ -733,7 +736,7 @@ def read_master_func(logfilename='', startB = dfAll.iloc[0]['beam'] while (r < threadCnt) and (n < rowCnt): - if (dfAll.loc[n]['beam']) != startB: + if (dfAll.iloc[n]['beam']) != startB: n+=1 else: rowsToProc.append((c, n)) @@ -744,7 +747,7 @@ def read_master_func(logfilename='', del c, r, n, startB, rowCnt # Fix no data in parallel - r = Parallel(n_jobs=threadCnt)(delayed(son._fixNoDat)(dfAll[r[0]:r[1]].copy().reset_index(drop=True), beams) for r in tqdm(rowsToProc)) + r = Parallel(n_jobs=safe_n_jobs(len(rowsToProc), threadCnt))(delayed(son._fixNoDat)(dfAll[r[0]:r[1]].copy().reset_index(drop=True), beams) for r in tqdm(rowsToProc)) gc.collect() # Concatenate results from parallel processing @@ -758,7 +761,7 @@ def read_master_func(logfilename='', # Slice dfAll by beam, update chunk_id, then save to file. for son in sonObjs: - df = dfAll[dfAll['beam'] == son.beam] + df = dfAll[dfAll['beam'] == son.beam].copy() if (len(df)%nchunk) != 0: rdr = nchunk-(len(df)%nchunk) @@ -783,7 +786,7 @@ def read_master_func(logfilename='', if len(lastChunk) <= (nchunk/2): df.loc[df['chunk_id']==c, 'chunk_id'] = c-1 - df.drop(columns = ['beam'], inplace=True) + df = df.drop(columns = ['beam']) # Check that last chunk has index anywhere in the chunk. ## If not, a bunch of NoData was added to the end. @@ -940,6 +943,13 @@ def read_master_func(logfilename='', df0 = df0[df0['filter'] == True] df1 = df1[df1['filter'] == True] + if df0.empty or df1.empty: + raise ValueError( + '\n\nFiltering removed all side-scan pings. No metadata remains to process. '\ + 'Adjust filtering parameters (max_heading_deviation, min_speed, max_speed, aoi, time_table) '\ + 'or reduce nchunk.' + ) + # Reasign the chunks df0 = son0._reassignChunks(df0) df1['chunk_id'] = df0['chunk_id'] @@ -1010,6 +1020,14 @@ def read_master_func(logfilename='', del son chunks = np.unique(chunks).astype(int) + + if len(chunks) == 0: + raise ValueError( + '\n\nNo valid side-scan chunks available for depth processing. '\ + 'This usually means prior filtering produced empty metadata CSVs. '\ + 'Relax filtering settings or reprocess from raw files.' + ) + # # Automatically estimate depth if detectDep > 0: # Check if depth detection dependencies are available @@ -1044,7 +1062,7 @@ def read_master_func(logfilename='', print('\n\tUsing binary thresholding...') # Parallel estimate depth for each chunk using appropriate method - r = Parallel(n_jobs=np.min([len(chunks), threadCnt]))(delayed(psObj._detectDepth)(detectDep, int(chunk), USE_GPU, tileFile) for chunk in tqdm(chunks)) + r = Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(psObj._detectDepth)(detectDep, int(chunk), USE_GPU, tileFile) for chunk in tqdm(chunks)) # store the depth predictions in the class for ret in r: @@ -1133,7 +1151,7 @@ def read_master_func(logfilename='', start_time = time.time() print("\n\nExporting bedpick plots to {}...".format(tileFile)) - Parallel(n_jobs=np.min([len(chunks), threadCnt]))(delayed(psObj._plotBedPick)(int(chunk), True, autoBed, tileFile) for chunk in tqdm(chunks)) + Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(psObj._plotBedPick)(int(chunk), True, autoBed, tileFile) for chunk in tqdm(chunks)) print("\nDone!") print("Time (s):", round(time.time() - start_time, ndigits=1)) @@ -1219,7 +1237,7 @@ def read_master_func(logfilename='', psObj.port.shadow = defaultdict() psObj.star.shadow = defaultdict() - r = Parallel(n_jobs=np.min([len(chunks), threadCnt]))(delayed(psObj._detectShadow)(remShadow, int(chunk), USE_GPU, False, tileFile) for chunk in tqdm(chunks)) + r = Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(psObj._detectShadow)(remShadow, int(chunk), USE_GPU, False, tileFile) for chunk in tqdm(chunks)) for ret in r: psObj.port.shadow[ret[0]] = ret[1] @@ -1274,7 +1292,7 @@ def read_master_func(logfilename='', # Calculate range-wise mean intensity for each chunk print('\n\tCalculating range-wise mean intensity for each chunk...') - chunk_means = Parallel(n_jobs= np.min([len(chunks), threadCnt]))(delayed(son._egnCalcChunkMeans)(i) for i in tqdm(chunks)) + chunk_means = Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._egnCalcChunkMeans)(i) for i in tqdm(chunks)) # Calculate global means print('\n\tCalculating range-wise global means...') @@ -1283,7 +1301,7 @@ def read_master_func(logfilename='', # Calculate egn min and max for each chunk print('\n\tCalculating EGN min and max values for each chunk...') - min_max = Parallel(n_jobs= np.min([len(chunks)]))(delayed(son._egnCalcMinMax)(i) for i in tqdm(chunks)) + min_max = Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._egnCalcMinMax)(i) for i in tqdm(chunks)) # Calculate global min max for each channel son._egnCalcGlobalMinMax(min_max) @@ -1335,7 +1353,7 @@ def read_master_func(logfilename='', chunks = chunks[:-1] # remove last chunk print('\n\tCalculating EGN corrected histogram for', son.beamName) - hist = Parallel(n_jobs= np.min([len(chunks), threadCnt]))(delayed(son._egnCalcHist)(i) for i in tqdm(chunks)) + hist = Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._egnCalcHist)(i) for i in tqdm(chunks)) print('\n\tCalculating global EGN corrected histogram') son._egnCalcGlobalHist(hist) @@ -1462,7 +1480,7 @@ def read_master_func(logfilename='', # Load sonMetaDF son._loadSonMeta() - Parallel(n_jobs= np.min([len(chunks), threadCnt]))(delayed(son._exportTilesSpd)(i, tileFile=imgType, spdCor=spdCor, mask_shdw=mask_shdw, maxCrop=maxCrop) for i in tqdm(chunks)) + Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._exportTilesSpd)(i, tileFile=imgType, spdCor=spdCor, mask_shdw=mask_shdw, maxCrop=maxCrop) for i in tqdm(chunks)) # for i in tqdm(chunks): # son._exportTilesSpd(i, tileFile=imgType, spdCor=spdCor, mask_shdw=mask_shdw, maxCrop=maxCrop) # sys.exit() diff --git a/pingmapper/main_rectify.py b/pingmapper/main_rectify.py index 8a2b224..4e4f433 100644 --- a/pingmapper/main_rectify.py +++ b/pingmapper/main_rectify.py @@ -357,7 +357,7 @@ def rectify_master_func(logfilename='', print('\n\tExporting', len(chunks), 'GeoTiffs for', son.beamName) # Parallel(n_jobs= np.min([len(sDF), threadCnt]))(delayed(son._rectSonHeadingMain)(sonarCoordsDF[sonarCoordsDF['chunk_id']==chunk], chunk) for chunk in tqdm(range(len(chunks)))) - Parallel(n_jobs= np.min([len(sDF), threadCnt]))(delayed(son._rectSonHeadingMain)(sDF[sDF['chunk_id']==chunk], chunk, heading=heading, interp_dist=rectInterpDist) for chunk in tqdm(chunks)) + Parallel(n_jobs=safe_n_jobs(len(sDF), threadCnt))(delayed(son._rectSonHeadingMain)(sDF[sDF['chunk_id']==chunk], chunk, heading=heading, interp_dist=rectInterpDist) for chunk in tqdm(chunks)) # for i in chunks: # # son._rectSonHeading(sonarCoordsDF[sonarCoordsDF['chunk_id']==i], i) # r = son._rectSonHeadingMain(sDF[sDF['chunk_id']==i], i, heading=heading, interp_dist=rectInterpDist) @@ -415,7 +415,7 @@ def rectify_master_func(logfilename='', # for i in chunks: # son._rectSonRubber(i, filter, cog, wgs=False) # sys.exit() - Parallel(n_jobs= np.min([len(chunks), threadCnt]))(delayed(son._rectSonRubber)(i, filter, cog, wgs=False) for i in tqdm(chunks)) + Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt))(delayed(son._rectSonRubber)(i, filter, cog, wgs=False) for i in tqdm(chunks)) son._cleanup() gc.collect() printUsage() diff --git a/pingmapper/utils/DRAFT_Workflows/avg_predictions_Mussel_WBL.py b/pingmapper/utils/DRAFT_Workflows/avg_predictions_Mussel_WBL.py index a1a93c3..bff21cd 100644 --- a/pingmapper/utils/DRAFT_Workflows/avg_predictions_Mussel_WBL.py +++ b/pingmapper/utils/DRAFT_Workflows/avg_predictions_Mussel_WBL.py @@ -234,7 +234,7 @@ def doWork(i, projDir, outDir_parent): os.mkdir(outDir_npz) - Parallel(n_jobs= np.min([len(npz_egn), threadCnt]), verbose=10)(delayed(npzAvg)(k, v, npz_raw, outDir_npz, son.nchunk) for k, v in npz_egn.items()) + Parallel(n_jobs=max(1, min(len(npz_egn), threadCnt)), verbose=10)(delayed(npzAvg)(k, v, npz_raw, outDir_npz, son.nchunk) for k, v in npz_egn.items()) ########################## @@ -265,7 +265,7 @@ def doWork(i, projDir, outDir_parent): # print('\n\n\n\n****Chunk', c) # son._pltSubClass('max', c, f, spdCor=1, maxCrop=True) - Parallel(n_jobs=np.min([len(toMap), threadCnt]), verbose=10)(delayed(son._pltSubClass)('max', c, f, spdCor=1, maxCrop=True, probs=probs) for c, f in toMap.items()) + Parallel(n_jobs=max(1, min(len(toMap), threadCnt)), verbose=10)(delayed(son._pltSubClass)('max', c, f, spdCor=1, maxCrop=True, probs=probs) for c, f in toMap.items()) del toMap @@ -313,7 +313,7 @@ def doWork(i, projDir, outDir_parent): # print(c, f) # psObj._mapSubstrate('max', c, f) - Parallel(n_jobs=np.min([len(toMap), threadCnt]), verbose=10)(delayed(psObj._mapSubstrate)('max', c, f) for c, f in toMap.items()) + Parallel(n_jobs=max(1, min(len(toMap), threadCnt)), verbose=10)(delayed(psObj._mapSubstrate)('max', c, f) for c, f in toMap.items()) del psObj diff --git a/pingmapper/utils/RawEGN_avg_predictions.py b/pingmapper/utils/RawEGN_avg_predictions.py index ea026bf..a4b9a8c 100644 --- a/pingmapper/utils/RawEGN_avg_predictions.py +++ b/pingmapper/utils/RawEGN_avg_predictions.py @@ -224,7 +224,7 @@ def doWork(i, projDir, outDir_parent): os.mkdir(outDir_npz) - Parallel(n_jobs= np.min([len(npz_egn), threadCnt]), verbose=10)(delayed(npzAvg)(k, v, npz_raw, outDir_npz, son.nchunk) for k, v in npz_egn.items()) + Parallel(n_jobs=max(1, min(len(npz_egn), threadCnt)), verbose=10)(delayed(npzAvg)(k, v, npz_raw, outDir_npz, son.nchunk) for k, v in npz_egn.items()) ########################## @@ -255,7 +255,7 @@ def doWork(i, projDir, outDir_parent): # print('\n\n\n\n****Chunk', c) # son._pltSubClass('max', c, f, spdCor=1, maxCrop=True) - Parallel(n_jobs=np.min([len(toMap), threadCnt]), verbose=10)(delayed(son._pltSubClass)('max', c, f, spdCor=1, maxCrop=True, probs=probs) for c, f in toMap.items()) + Parallel(n_jobs=max(1, min(len(toMap), threadCnt)), verbose=10)(delayed(son._pltSubClass)('max', c, f, spdCor=1, maxCrop=True, probs=probs) for c, f in toMap.items()) del toMap @@ -303,7 +303,7 @@ def doWork(i, projDir, outDir_parent): # print(c, f) # psObj._mapSubstrate('max', c, f) - Parallel(n_jobs=np.min([len(toMap), threadCnt]), verbose=10)(delayed(psObj._mapSubstrate)('max', c, f) for c, f in toMap.items()) + Parallel(n_jobs=max(1, min(len(toMap), threadCnt)), verbose=10)(delayed(psObj._mapSubstrate)('max', c, f) for c, f in toMap.items()) del psObj diff --git a/pingmapper/utils/Substrate_Summaries/02_gen_summary_stamp_shps.py b/pingmapper/utils/Substrate_Summaries/02_gen_summary_stamp_shps.py index 13ae184..1ed9173 100644 --- a/pingmapper/utils/Substrate_Summaries/02_gen_summary_stamp_shps.py +++ b/pingmapper/utils/Substrate_Summaries/02_gen_summary_stamp_shps.py @@ -259,7 +259,7 @@ def getBearingLine(x, y, rr, rl, d): # Try splitting polygon with lines # Get polygon to split - to_split = covShp.loc[0].geometry + to_split = covShp.iloc[0].geometry # Split with first line polys = split(to_split, beginLine) @@ -277,7 +277,7 @@ def getBearingLine(x, y, rr, rl, d): splitGDF = splitGDF.dissolve(by=None) # Get polygon to split - to_split = splitGDF.loc[0].geometry + to_split = splitGDF.iloc[0].geometry # Split with second line polys = split(to_split, endLine) diff --git a/pingmapper/utils/Substrate_Summaries/summarize_project_substrate.py b/pingmapper/utils/Substrate_Summaries/summarize_project_substrate.py index aeec63e..559dcce 100644 --- a/pingmapper/utils/Substrate_Summaries/summarize_project_substrate.py +++ b/pingmapper/utils/Substrate_Summaries/summarize_project_substrate.py @@ -120,7 +120,7 @@ def extendLines(pnts, d): def calcSinuosity(df): # Calculate channel length - channel_len = df.iloc[-1]['trk_dist'] - df.loc[0]['trk_dist'] + channel_len = df.iloc[-1]['trk_dist'] - df.iloc[0]['trk_dist'] # Calculate downvalley length x_1 = df.iloc[0]['trk_utm_es'] @@ -352,7 +352,7 @@ def doWork(i, projDir): # Try splitting polygon with lines # Get polygon to split - to_split = diss.loc[0].geometry + to_split = diss.iloc[0].geometry # Split with first line polys = split(to_split, beginLine) @@ -387,7 +387,7 @@ def doWork(i, projDir): # Get polygon to split - to_split = splitGDF.loc[0].geometry + to_split = splitGDF.iloc[0].geometry # Split with second line polys = split(to_split, endLine) @@ -583,7 +583,7 @@ def doWork(i, projDir): proj_cnt = len(projDirs) -Parallel(n_jobs= np.min([len(projDirs), threadCnt]), verbose=10)(delayed(doWork)(i, p) for i, p in enumerate(projDirs)) +Parallel(n_jobs=safe_n_jobs(len(projDirs), threadCnt), verbose=10)(delayed(doWork)(i, p) for i, p in enumerate(projDirs)) ## For testing ##projDirs = projDirs[:10] diff --git a/pingmapper/utils/main_mosaic_transects.py b/pingmapper/utils/main_mosaic_transects.py index 0849eae..4a8fd5d 100644 --- a/pingmapper/utils/main_mosaic_transects.py +++ b/pingmapper/utils/main_mosaic_transects.py @@ -309,7 +309,7 @@ def average(in_ar, out_ar, xoff, yoff, xsize, ysize, raster_xsize,raster_ysize, # Smooth tracklines print("\nSmoothing tracklines...") -Parallel(n_jobs= np.min([len(portstar), threadCnt]), verbose=10)(delayed(smoothTrackline)(sons) for sons in portstar) +Parallel(n_jobs=safe_n_jobs(len(portstar), threadCnt), verbose=10)(delayed(smoothTrackline)(sons) for sons in portstar) # Store smthTrkFile in rectObj for sons in portstar: @@ -318,7 +318,7 @@ def average(in_ar, out_ar, xoff, yoff, xsize, ysize, raster_xsize,raster_ysize, # Calculate range extent coordinates print("\nCalculating, smoothing, and interpolating range extent coordinates...") -Parallel(n_jobs= np.min([len(portstar), threadCnt]), verbose=10)(delayed(rangeCoordinates)(sons) for sons in portstar) +Parallel(n_jobs=safe_n_jobs(len(portstar), threadCnt), verbose=10)(delayed(rangeCoordinates)(sons) for sons in portstar) for sons in portstar: for son in sons: @@ -599,7 +599,7 @@ def average(in_ar, out_ar, xoff, yoff, xsize, ysize, raster_xsize,raster_ysize, son._getSonColorMap(son_colorMap) print('\n\tExporting', len(chunks), 'GeoTiffs for', son.beamName) - Parallel(n_jobs= np.min([len(chunks), threadCnt]), verbose=10)(delayed(son._rectSonParallel)(i, filter, wgs=False) for i in chunks) + Parallel(n_jobs=safe_n_jobs(len(chunks), threadCnt), verbose=10)(delayed(son._rectSonParallel)(i, filter, wgs=False) for i in chunks) del son.sonMetaDF try: