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
2 changes: 1 addition & 1 deletion app/routers/hydrofabric/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -255,7 +255,7 @@ def _release_sem() -> None:
spatial_names: list[str] = []
nonspatial_names: list[str] = []
for name, data in output_layers.items():
if isinstance(data, gpd.GeoDataFrame) and len(data) > 0:
if isinstance(data, gpd.GeoDataFrame):
spatial_names.append(name)
elif not isinstance(data, gpd.GeoDataFrame):
nonspatial_names.append(name)
Expand Down
24 changes: 8 additions & 16 deletions src/icefabric/hydrofabric/subset_nhf.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,9 @@ def _build_upstream_dict_from_nexus(
"""Build upstream connectivity dictionary from flowpath nexus connections."""
fp_pl = flowpaths_pl.with_columns(
[
pl.col(edge_id).cast(pl.Int32),
pl.col(f"up_{node_id}").cast(pl.Int32),
pl.col(f"dn_{node_id}").cast(pl.Int32),
pl.col(edge_id).cast(pl.Int64),
pl.col(f"up_{node_id}").cast(pl.Int64),
pl.col(f"dn_{node_id}").cast(pl.Int64),
]
)
nexus_to_downstream = fp_pl.select(
Expand Down Expand Up @@ -253,15 +253,13 @@ def generate_subset_from_ids(
f = {
"fp": ex.submit(source.load_filtered, "flowpaths", "fp_id", flowpath_ids),
"div": ex.submit(source.load_filtered, "divides", "div_id", flowpath_ids),
"wb": ex.submit(source.load_filtered, "waterbodies", "fp_id", flowpath_ids),
"gages": ex.submit(source.load_filtered, "gages", "fp_id", flowpath_ids),
"ref_fp": ex.submit(source.load_filtered, "reference_flowpaths", "div_id", flowpath_ids),
"lakes": ex.submit(source.load_filtered, "lakes", "fp_id", flowpath_ids),
"nhd": ex.submit(source.load_filtered, "nhd", "ref_id", flowpath_ids),
}
subset_fp = f["fp"].result()
subset_div = f["div"].result()
subset_wb = f["wb"].result()
subset_gages = f["gages"].result()
subset_ref_fp = f["ref_fp"].result()
subset_lakes = f["lakes"].result()
Expand Down Expand Up @@ -292,9 +290,9 @@ def generate_subset_from_ids(
+ subset_fp.filter(pl.col("dn_nex_id").is_not_null())["dn_nex_id"].cast(pl.Int64).to_list()
)
all_v_fp_ids = set(subset_ref_fp["virtual_fp_id"].to_list())
wb_hy_ids = subset_wb["hy_id"].to_list() if "hy_id" in subset_wb.columns else []
lakes_hy_ids = subset_lakes["hy_id"].to_list() if "hy_id" in subset_lakes.columns else []
gage_hy_ids = subset_gages["hy_id"].to_list() if "hy_id" in subset_gages.columns else []
all_hy_ids = set(wb_hy_ids + gage_hy_ids)
all_hy_ids = set(lakes_hy_ids + gage_hy_ids)

# Wave 2: nex_id/virtual_fp_id filtered (parallel)
with ThreadPoolExecutor(max_workers=2) as ex:
Expand Down Expand Up @@ -345,7 +343,6 @@ def generate_subset_from_ids(
"divides": pl_to_gdf(subset_div, crs=crs),
"virtual_nexus": pl_to_gdf(subset_v_nex, crs=crs),
"virtual_flowpaths": pl_to_gdf(subset_v_fp, crs=crs),
"waterbodies": pl_to_gdf(subset_wb, crs=crs) if len(subset_wb) > 0 else subset_wb.to_pandas(),
"gages": pl_to_gdf(subset_gages, crs=crs) if len(subset_gages) > 0 else subset_gages.to_pandas(),
"lakes": pl_to_gdf(subset_lakes, crs=crs) if len(subset_lakes) > 0 else subset_lakes.to_pandas(),
"reference_flowpaths": subset_ref_fp.to_pandas(),
Expand All @@ -363,7 +360,6 @@ def generate_subset_from_ids(
"divides",
"virtual_nexus",
"virtual_flowpaths",
"waterbodies",
"gages",
"lakes",
]
Expand Down Expand Up @@ -425,14 +421,12 @@ def generate_subset_virtual_only(
f = {
"fp": ex.submit(source.load_filtered, "flowpaths", "div_id", div_ids),
"div": ex.submit(source.load_filtered, "divides", "div_id", div_ids),
"wb": ex.submit(source.load_filtered, "waterbodies", "div_id", div_ids),
"gages": ex.submit(source.load_filtered, "gages", "div_id", div_ids),
"ref_fp": ex.submit(source.load_filtered, "reference_flowpaths", "div_id", div_ids),
"lakes": ex.submit(source.load_filtered, "lakes", "div_id", div_ids),
}
subset_fp = f["fp"].result()
subset_div = f["div"].result()
subset_wb = f["wb"].result()
subset_gages = f["gages"].result()
subset_ref_fp = f["ref_fp"].result()
subset_lakes = f["lakes"].result()
Expand Down Expand Up @@ -463,10 +457,10 @@ def generate_subset_virtual_only(
.to_list()
)

# Derive hydrolocation IDs from gages / waterbodies in this divide
wb_hy_ids = subset_wb["hy_id"].to_list() if "hy_id" in subset_wb.columns else []
# Derive hydrolocation IDs from gages and lakes in this divide
lakes_hy_ids = subset_lakes["hy_id"].to_list() if "hy_id" in subset_lakes.columns else []
gage_hy_ids = subset_gages["hy_id"].to_list() if "hy_id" in subset_gages.columns else []
all_hy_ids = set(wb_hy_ids + gage_hy_ids)
all_hy_ids = set(lakes_hy_ids + gage_hy_ids)

# Wave 2: nexus / virtual_nexus / hydrolocations
with ThreadPoolExecutor(max_workers=2) as ex:
Expand Down Expand Up @@ -503,7 +497,6 @@ def generate_subset_virtual_only(
"divides": pl_to_gdf(subset_div, crs=crs),
"virtual_nexus": pl_to_gdf(subset_v_nex, crs=crs),
"virtual_flowpaths": pl_to_gdf(subset_v_fp, crs=crs),
"waterbodies": (pl_to_gdf(subset_wb, crs=crs) if len(subset_wb) > 0 else subset_wb.to_pandas()),
"gages": (pl_to_gdf(subset_gages, crs=crs) if len(subset_gages) > 0 else subset_gages.to_pandas()),
"lakes": (pl_to_gdf(subset_lakes, crs=crs) if len(subset_lakes) > 0 else subset_lakes.to_pandas()),
"reference_flowpaths": subset_ref_fp.to_pandas(),
Expand All @@ -521,7 +514,6 @@ def generate_subset_virtual_only(
"divides",
"virtual_nexus",
"virtual_flowpaths",
"waterbodies",
"gages",
"lakes",
]
Expand Down
2 changes: 0 additions & 2 deletions src/icefabric/schemas/iceberg_tables/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,13 @@
ReferenceFlowpaths,
VirtualFlowpaths,
VirtualNexus,
Waterbodies,
)

nhf_layers = {
"divides": Divides,
"flowpaths": Flowpaths,
"nexus": Nexus,
"reference_flowpaths": ReferenceFlowpaths,
"waterbodies": Waterbodies,
"gages": Gages,
"virtual_flowpaths": VirtualFlowpaths,
"virtual_nexus": VirtualNexus,
Expand Down
42 changes: 41 additions & 1 deletion src/icefabric/schemas/iceberg_tables/hydrofabric_update.py
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,7 @@ def columns(cls) -> list[str]:
"expon",
"max_gw_storage",
"geometry",
"gid",
]

@classmethod
Expand Down Expand Up @@ -395,6 +396,7 @@ def schema(cls) -> Schema:
NestedField(68, "expon", DoubleType(), required=False, doc=desc[67]),
NestedField(69, "max_gw_storage", DoubleType(), required=False, doc=desc[68]),
NestedField(70, "geometry", BinaryType(), required=False, doc=desc[69]),
NestedField(71, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand Down Expand Up @@ -480,6 +482,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("expon", pa.float64(), nullable=True),
pa.field("max_gw_storage", pa.float64(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -596,6 +599,8 @@ def columns(cls) -> list[str]:
"r_ml",
"fp_to_id",
"geometry",
"gid",
"terminalpa",
]

@classmethod
Expand Down Expand Up @@ -673,6 +678,10 @@ def schema(cls) -> Schema:
NestedField(29, "r_ml", FloatType(), required=False, doc=desc[28]),
NestedField(30, "fp_to_id", LongType(), required=False, doc=desc[29]),
NestedField(31, "geometry", BinaryType(), required=False, doc=desc[30]),
NestedField(32, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
NestedField(
33, "terminalpa", LongType(), required=False, doc="Terminal path grouping identifier"
),
identifier_field_ids=[1],
)

Expand Down Expand Up @@ -719,6 +728,8 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("r_ml", pa.float32(), nullable=True),
pa.field("fp_to_id", pa.int64(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
pa.field("terminalpa", pa.int64(), nullable=True),
]
)

Expand Down Expand Up @@ -754,6 +765,7 @@ def columns(cls) -> list[str]:
"dn_fp_id",
"vpu_id",
"geometry",
"gid",
]

@classmethod
Expand All @@ -777,6 +789,7 @@ def schema(cls) -> Schema:
NestedField(2, "dn_fp_id", LongType(), required=False, doc=desc[1]),
NestedField(3, "vpu_id", StringType(), required=False, doc=desc[2]),
NestedField(4, "geometry", BinaryType(), required=False, doc=desc[3]),
NestedField(5, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand All @@ -796,6 +809,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("dn_fp_id", pa.int64(), nullable=True),
pa.field("vpu_id", pa.string(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -837,6 +851,7 @@ def columns(cls) -> list[str]:
"div_id",
"mainstem_virtual_fp_id",
"segment_order",
"gid",
]

@classmethod
Expand Down Expand Up @@ -864,6 +879,7 @@ def schema(cls) -> Schema:
NestedField(4, "div_id", LongType(), required=False, doc=desc[3]),
NestedField(5, "mainstem_virtual_fp_id", LongType(), required=False, doc=desc[4]),
NestedField(6, "segment_order", LongType(), required=False, doc=desc[5]),
NestedField(7, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand All @@ -885,6 +901,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("div_id", pa.int64(), nullable=True),
pa.field("mainstem_virtual_fp_id", pa.int64(), nullable=True),
pa.field("segment_order", pa.int64(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -974,6 +991,7 @@ def columns(cls) -> list[str]:
"dn_virtual_nex_id",
"virtual_fp_id",
"geometry",
"gid",
]

@classmethod
Expand Down Expand Up @@ -1070,6 +1088,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("dn_virtual_nex_id", pa.float64(), nullable=True),
pa.field("virtual_fp_id", pa.float64(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -1135,6 +1154,7 @@ def columns(cls) -> list[str]:
"mainstem_virtual_fp_id",
"segment_order",
"geometry",
"gid",
]

@classmethod
Expand Down Expand Up @@ -1178,6 +1198,7 @@ def schema(cls) -> Schema:
NestedField(12, "mainstem_virtual_fp_id", DoubleType(), required=False, doc=desc[11]),
NestedField(13, "segment_order", DoubleType(), required=False, doc=desc[12]),
NestedField(14, "geometry", BinaryType(), required=False, doc=desc[13]),
NestedField(15, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand Down Expand Up @@ -1207,6 +1228,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("mainstem_virtual_fp_id", pa.float64(), nullable=True),
pa.field("segment_order", pa.float64(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -1257,6 +1279,7 @@ def columns(cls) -> list[str]:
"percentage_area_contribution",
"vpu_id",
"geometry",
"gid",
]

@classmethod
Expand Down Expand Up @@ -1290,6 +1313,7 @@ def schema(cls) -> Schema:
NestedField(7, "percentage_area_contribution", DoubleType(), required=False, doc=desc[6]),
NestedField(8, "vpu_id", StringType(), required=False, doc=desc[7]),
NestedField(9, "geometry", BinaryType(), required=False, doc=desc[8]),
NestedField(10, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand All @@ -1314,6 +1338,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("percentage_area_contribution", pa.float64(), nullable=True),
pa.field("vpu_id", pa.string(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -1349,6 +1374,7 @@ def columns(cls) -> list[str]:
"dn_virtual_fp_id",
"vpu_id",
"geometry",
"gid",
]

@classmethod
Expand All @@ -1372,6 +1398,7 @@ def schema(cls) -> Schema:
NestedField(2, "dn_virtual_fp_id", LongType(), required=False, doc=desc[1]),
NestedField(3, "vpu_id", StringType(), required=False, doc=desc[2]),
NestedField(4, "geometry", BinaryType(), required=False, doc=desc[3]),
NestedField(5, "gid", StringType(), required=False, doc="Geolocation Plus Code identifier"),
identifier_field_ids=[1],
)

Expand All @@ -1391,6 +1418,7 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("dn_virtual_fp_id", pa.int64(), nullable=True),
pa.field("vpu_id", pa.string(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
pa.field("gid", pa.string(), nullable=True),
]
)

Expand Down Expand Up @@ -1451,6 +1479,10 @@ class Lakes:
Reservoir index for Medium Range configuration
reservoir_index_Short_Range : float
Reservoir index for Short Range configuration
dam_id : str
Dam identifier
nidid : str
National Inventory of Dams identifier
geometry : binary
Spatial Geometry (POINT format) - stored in WKB binary format
"""
Expand Down Expand Up @@ -1484,6 +1516,8 @@ def columns(cls) -> list[str]:
"reservoir_index_GDL_AK",
"reservoir_index_Medium_Range",
"reservoir_index_Short_Range",
"dam_id",
"nidid",
"geometry",
]

Expand Down Expand Up @@ -1516,6 +1550,8 @@ def schema(cls) -> Schema:
"Reservoir index for GDL AK configuration",
"Reservoir index for Medium Range configuration",
"Reservoir index for Short Range configuration",
"Dam identifier",
"National Inventory of Dams identifier",
"Spatial Geometry (POINT format) - stored in WKB binary format",
]
return Schema(
Expand Down Expand Up @@ -1544,7 +1580,9 @@ def schema(cls) -> Schema:
NestedField(23, "reservoir_index_GDL_AK", DoubleType(), required=False, doc=desc[22]),
NestedField(24, "reservoir_index_Medium_Range", DoubleType(), required=False, doc=desc[23]),
NestedField(25, "reservoir_index_Short_Range", DoubleType(), required=False, doc=desc[24]),
NestedField(26, "geometry", BinaryType(), required=False, doc=desc[25]),
NestedField(26, "dam_id", StringType(), required=False, doc=desc[25]),
NestedField(27, "nidid", StringType(), required=False, doc=desc[26]),
NestedField(28, "geometry", BinaryType(), required=False, doc=desc[27]),
identifier_field_ids=[1],
)

Expand Down Expand Up @@ -1578,6 +1616,8 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("reservoir_index_GDL_AK", pa.float64(), nullable=True),
pa.field("reservoir_index_Medium_Range", pa.float64(), nullable=True),
pa.field("reservoir_index_Short_Range", pa.float64(), nullable=True),
pa.field("dam_id", pa.string(), nullable=True),
pa.field("nidid", pa.string(), nullable=True),
pa.field("geometry", pa.binary(), nullable=True),
]
)
Expand Down
13 changes: 5 additions & 8 deletions src/icefabric/schemas/iceberg_tables/nhf_snapshots.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,6 @@ def columns(cls) -> list[str]:
"flowpaths",
"nexus",
"reference_flowpaths",
"waterbodies",
"gages",
"virtual_flowpaths",
"virtual_nexus",
Expand All @@ -57,12 +56,11 @@ def schema(cls) -> Schema:
NestedField(2, "flowpaths", LongType(), required=False),
NestedField(3, "nexus", LongType(), required=False),
NestedField(4, "reference_flowpaths", LongType(), required=False),
NestedField(5, "waterbodies", LongType(), required=False),
NestedField(6, "gages", LongType(), required=False),
NestedField(7, "virtual_flowpaths", LongType(), required=False),
NestedField(8, "virtual_nexus", LongType(), required=False),
NestedField(9, "hydrolocations", LongType(), required=False),
NestedField(10, "lakes", LongType(), required=False),
NestedField(5, "gages", LongType(), required=False),
NestedField(6, "virtual_flowpaths", LongType(), required=False),
NestedField(7, "virtual_nexus", LongType(), required=False),
NestedField(8, "hydrolocations", LongType(), required=False),
NestedField(9, "lakes", LongType(), required=False),
)

@classmethod
Expand All @@ -80,7 +78,6 @@ def arrow_schema(cls) -> pa.Schema:
pa.field("flowpaths", pa.int64(), nullable=True),
pa.field("nexus", pa.int64(), nullable=True),
pa.field("reference_flowpaths", pa.int64(), nullable=True),
pa.field("waterbodies", pa.int64(), nullable=True),
pa.field("gages", pa.int64(), nullable=True),
pa.field("virtual_flowpaths", pa.int64(), nullable=True),
pa.field("virtual_nexus", pa.int64(), nullable=True),
Expand Down
Loading