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
214 changes: 169 additions & 45 deletions ati/mod0Helper/avaDirectory/avaDirBuildFromFlowPy.py

Large diffs are not rendered by default.

71 changes: 44 additions & 27 deletions ati/mod0Helper/avaDirectory/avaDirResults.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,6 @@
import geopandas as gpd

import ati.mod0Helper.dataUtils as dataUtils
from ati.mod0Helper.dataUtils import relPath

log = logging.getLogger(__name__)
logging.getLogger("pyogrio").setLevel(logging.WARNING)
Expand All @@ -74,22 +73,18 @@ def runAvaDirResults(cfg, workFlowDir):
log.info("Step 15: Start AvaDirectory Results build...")

avaCfg = cfg["avaDIRECTORY"]
main = cfg["MAIN"]

# --- Resolve core directories dynamically ---
rootDir = Path(main["workDir"]) / main["project"] / main["ID"]

# Source: FlowPy output rasters (raw data)
avaDirData = rootDir / "11_avaDirectoryData"
avaDirData = Path(workFlowDir["avaDirDir"])
# Destination: library of merged outputs
avaDirLib = rootDir / "12_avaDirectory"
avaDirLib = Path(workFlowDir["avaDirResultsDir"])
# Ensure destination exists
avaDirLib.mkdir(parents=True, exist_ok=True)

cairosDir = Path(workFlowDir["cairosDir"])
log.info("Step 15: Using AvaDirectoryData = %s", relPath(avaDirData, cairosDir))
log.info("Step 15: Using AvaDirectoryData = %s", dataUtils.relPath(avaDirData, cairosDir))
log.info(
"Step 15: Writing outputs to AvaDirectory = %s", relPath(avaDirLib, cairosDir)
"Step 15: Writing outputs to AvaDirectory = %s", dataUtils.relPath(avaDirLib, cairosDir)
)

# --- Check that input AvaDirectoryType exists ---
Expand All @@ -98,7 +93,7 @@ def runAvaDirResults(cfg, workFlowDir):
if not avaTypeParquet.exists() and not avaTypeGeoJSON.exists():
log.error(
"Step 15: No avaDirectoryType.* found in %s",
relPath(avaDirLib, cairosDir),
dataUtils.relPath(avaDirLib, cairosDir),
)
return

Expand All @@ -125,7 +120,7 @@ def runAvaDirResults(cfg, workFlowDir):

# --- Raster filename patterns ---
typePatterns = {
"inputPRA": "-area_m.tif",
"inputPRA": ("-praAreaM.tif", "-area_m.tif"),
"cellCounts": "_cellCounts_lzw.tif",
"zDelta": "_zdelta_lzw.tif",
"zDelta_sized": "_zdelta_sized_lzw.tif",
Expand Down Expand Up @@ -209,8 +204,10 @@ def _buildFileIndex(avaDirData: Path, typePatterns: dict) -> dict:
pra = int(praStr)
except ValueError:
continue
for t, pat in typePatterns.items():
if pat in fname:
for t, patterns in typePatterns.items():
if isinstance(patterns, str):
patterns = (patterns,)
if any(pattern in fname for pattern in patterns):
index.setdefault((pra, rid), {})[t] = str(tifPath.resolve())
break
return index
Expand All @@ -222,16 +219,26 @@ def _loadOrBuildFileIndex(avaDirData, indexFile, typePatterns, forceRebuild, cai
try:
with open(indexFile, "rb") as f:
index = pickle.load(f)
log.info(
"Step 15: Loaded cached index (%d entries) from %s",
len(index),
relPath(indexFile, cairosDir),
indexedInputPra = any("inputPRA" in entry for entry in index.values())
inputPraPatterns = typePatterns["inputPRA"]
availableInputPra = any(
any(com4Dir.glob(f"*{pattern}"))
for com4Dir in Path(avaDirData).glob("com4_*")
for pattern in inputPraPatterns
)
return index
if availableInputPra and not indexedInputPra:
log.info("Step 15: Cached index has no input PRA paths; rebuilding index.")
else:
log.info(
"Step 15: Loaded cached index (%d entries) from %s",
len(index),
dataUtils.relPath(indexFile, cairosDir),
)
return index
except Exception as e:
log.warning("Step 15: Failed to load cached index (%s), rebuilding...", e)

log.info("Step 15: Building file index from %s", relPath(avaDirData, cairosDir))
log.info("Step 15: Building file index from %s", dataUtils.relPath(avaDirData, cairosDir))
index = _buildFileIndex(avaDirData, typePatterns)
with open(indexFile, "wb") as f:
pickle.dump(index, f)
Expand Down Expand Up @@ -261,11 +268,21 @@ def _makeAvaDirResults(
if outParquet.exists() and not forceRebuild and writeParquet:
try:
avaDir = gpd.read_parquet(outParquet)
log.info(
"Step 15: Loaded cached AvaDirectoryResults (%d features)",
len(avaDir),
indexedInputPra = any("inputPRA" in entry for entry in fileIndex.values())
cachedInputPra = (
"pathInputpra" in avaDir.columns
and avaDir["pathInputpra"].notna().any()
)
return avaDir
if indexedInputPra and not cachedInputPra:
log.info(
"Step 15: Cached results have no input PRA paths; rebuilding results."
)
else:
log.info(
"Step 15: Loaded cached AvaDirectoryResults (%d features)",
len(avaDir),
)
return avaDir
except Exception as e:
log.warning("Step 15: Failed to load cached Results (%s), rebuilding...", e)

Expand Down Expand Up @@ -332,7 +349,7 @@ def _rel(p):
avaDir.drop(columns="geometry", errors="ignore").to_csv(outCsv, index=False)
log.info(
"Step 15: Wrote CSV AvaDirectoryResults to %s",
relPath(outCsv, cairosDir),
dataUtils.relPath(outCsv, cairosDir),
)
except Exception as e:
log.warning("Step 15: CSV write warning: %s", e)
Expand All @@ -342,7 +359,7 @@ def _rel(p):
avaDir.to_file(outGeoJson, driver="GeoJSON")
log.info(
"Step 15: Wrote GeoJSON AvaDirectoryResults to %s",
relPath(outGeoJson, cairosDir),
dataUtils.relPath(outGeoJson, cairosDir),
)
except Exception as e:
log.warning("Step 15: GeoJSON write warning: %s", e)
Expand All @@ -352,14 +369,14 @@ def _rel(p):
avaDir.to_parquet(outParquet, index=False)
log.info(
"Step 15: Wrote Parquet AvaDirectoryResults to %s",
relPath(outParquet, cairosDir),
dataUtils.relPath(outParquet, cairosDir),
)
except Exception as e:
log.warning("Step 15: Parquet write warning: %s", e)

log.info(
"Step 15: Wrote AvaDirectoryResults (%d features) to %s",
len(avaDir),
relPath(avaDirLib, cairosDir),
dataUtils.relPath(avaDirLib, cairosDir),
)
return avaDir
2 changes: 1 addition & 1 deletion ati/mod0Helper/avaDirectory/avaDirResultsStats.py
Original file line number Diff line number Diff line change
Expand Up @@ -583,7 +583,7 @@ def _dupSubsetColsScenarioIndependent(df: pd.DataFrame, modType: str) -> List[st
"""
base = [
"praID", "modType",
"LKGebiet", "LKGebietID", "LKRegion", "LWDGebietID",
"LKGebiet", "LKGebietID", "LKRegion",
"praAreaM", "praAreaSized", "praAreaVol",
"praElevBand", "praElevBandRule",
"praElevMax", "praElevMean", "praElevMin",
Expand Down
29 changes: 14 additions & 15 deletions ati/mod0Helper/avaDirectory/avaDirType.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,6 @@
import pandas as pd

import ati.mod0Helper.dataUtils as dataUtils
from ati.mod0Helper.dataUtils import relPath

import sys
from functools import partial
Expand Down Expand Up @@ -76,14 +75,12 @@ def runAvaDirType(cfg, workFlowDir):
log.info("Step 14: Start AvaDirectory Type build...")

wf = cfg["WORKFLOW"]
main = cfg["MAIN"]
avaCfg = cfg["avaDIRECTORY"]

# --- Resolve directories dynamically ---
rootDir = Path(main["workDir"]) / main["project"] / main["ID"]
# --- Resolve directories supplied by the active workflow ---
cairosDir = Path(workFlowDir["cairosDir"])

avaDirLib = rootDir / "12_avaDirectory"
avaDirLib = Path(workFlowDir["avaDirTypeDir"])
avaDirLib.mkdir(parents=True, exist_ok=True)

# --- Mode flags ---
Expand Down Expand Up @@ -121,15 +118,17 @@ def runAvaDirType(cfg, workFlowDir):
if readScenarioParquet:
# Scenario parquet files are written by Step 13 into FlowPy folder tree:
# 09_flowPyBigDataStructure/pra*/Size*/{dry,wet}/Map/singleAvaDir/com4_*/avaScenario.parquet
flowPyRoot = rootDir / "09_flowPyBigDataStructure"
flowPySourceDir = workFlowDir.get("flowPySourceDir") or workFlowDir["flowPyRunDir"]
flowPyRoot = Path(flowPySourceDir)
if not flowPyRoot.exists():
log.error(
"Step 14: FlowPy root missing: %s", relPath(flowPyRoot, cairosDir)
"Step 14: FlowPy root missing: %s", dataUtils.relPath(flowPyRoot, cairosDir)
)
return

pattern = str(
flowPyRoot
/ "**"
/ "pra*"
/ "Size*"
/ "*"
Expand All @@ -138,23 +137,23 @@ def runAvaDirType(cfg, workFlowDir):
/ "com4_*"
/ "avaScenLeaf_com4_*.parquet"
)
inputFiles = sorted(glob.glob(pattern))
inputFiles = sorted(glob.glob(pattern, recursive=True))
log.info("Step 14: Found %d scenario parquet files", len(inputFiles))

elif readSingleAvaGeoJSON:
# Legacy mode: read from 11_avaDirectoryData/com4_*/praID*.geojson
avaDirData = rootDir / "11_avaDirectoryData"
avaDirData = Path(workFlowDir["avaDirDir"])
if not avaDirData.exists():
log.warning(
"Step 14: Expected AvaDirectoryData missing: %s",
relPath(avaDirData, cairosDir),
dataUtils.relPath(avaDirData, cairosDir),
)
return

com4Folders = sorted(glob.glob(str(avaDirData / "com4_*")))
if not com4Folders:
log.warning(
"Step 14: No com4_* folders found in %s", relPath(avaDirData, cairosDir)
"Step 14: No com4_* folders found in %s", dataUtils.relPath(avaDirData, cairosDir)
)
return

Expand Down Expand Up @@ -185,7 +184,7 @@ def runAvaDirType(cfg, workFlowDir):
gdf = dataUtils.readGeoData(fp)
allChunks.append(gdf)
except Exception:
log.exception("Step 14: Failed to read %s", relPath(Path(fp), cairosDir))
log.exception("Step 14: Failed to read %s", dataUtils.relPath(Path(fp), cairosDir))

if not allChunks:
log.warning("Step 14: All reads failed → no output created.")
Expand Down Expand Up @@ -246,21 +245,21 @@ def runAvaDirType(cfg, workFlowDir):
merged.drop(columns="geometry", errors="ignore").to_csv(
csvPath, index=False
)
log.info("Step 14: Wrote CSV to %s", relPath(csvPath, cairosDir))
log.info("Step 14: Wrote CSV to %s", dataUtils.relPath(csvPath, cairosDir))
except Exception as e:
log.warning("Step 14: CSV write warning: %s", e)

if writeGeoJSON:
try:
dataUtils.writeGeoData(merged, geojsonPath)
log.info("Step 14: Wrote GeoJSON to %s", relPath(geojsonPath, cairosDir))
log.info("Step 14: Wrote GeoJSON to %s", dataUtils.relPath(geojsonPath, cairosDir))
except Exception as e:
log.warning("Step 14: GeoJSON write warning: %s", e)

if writeParquet:
try:
merged.to_parquet(parquetPath, index=False)
log.info("Step 14: Wrote Parquet to %s", relPath(parquetPath, cairosDir))
log.info("Step 14: Wrote Parquet to %s", dataUtils.relPath(parquetPath, cairosDir))
except Exception as e:
log.warning("Step 14: Parquet write warning: %s", e)

Expand Down
Loading