Files
pbf/navsea_detection_pipeline.py
2026-03-17 19:48:15 +08:00

435 lines
16 KiB
Python

from __future__ import annotations
import os
from dataclasses import dataclass
import pymysql
@dataclass(frozen=True)
class DbConfig:
host: str = os.getenv("NAVSEA_DB_HOST", "localhost")
port: int = int(os.getenv("NAVSEA_DB_PORT", "3306"))
user: str = os.getenv("NAVSEA_DB_USER", "root")
password: str = os.getenv("NAVSEA_DB_PASSWORD", "2chi9ks2")
database: str = os.getenv("NAVSEA_DB_NAME", "pbf_analysis")
unix_socket: str | None = os.getenv("NAVSEA_DB_SOCKET", "/tmp/mysql.sock")
FAMILY_DEFAULT_CAPABILITIES = {
"hazard": {"collision"},
"navigation_aid": {"navigation_mark"},
"bathymetry": {"depth_reference"},
"seabed": {"bottom_reference"},
"fishery": {"entangle"},
"route": {"route_reference"},
"boundary": {"boundary_reference"},
"anchorage": {"anchorage_reference"},
"restricted_area": {"boundary_reference"},
"infrastructure": {"structure_reference"},
"utility": {"utility_hazard"},
"landmark": {"landmark_reference"},
"place": {"place_reference"},
"water_area": {"water_reference"},
"monitoring": {"traffic_monitoring"},
"overlay_support": {"render_support"},
"unknown": {"review_required"},
}
class DetectionPipeline:
def __init__(self, config: DbConfig) -> None:
self.config = config
def connect(self):
kwargs = {
"host": self.config.host,
"port": self.config.port,
"user": self.config.user,
"password": self.config.password,
"database": self.config.database,
"charset": "utf8mb4",
"autocommit": False,
}
if self.config.unix_socket and self.config.host in {"localhost", "127.0.0.1"}:
kwargs["unix_socket"] = self.config.unix_socket
return pymysql.connect(**kwargs)
def run(self) -> None:
with self.connect() as conn:
with conn.cursor() as cur:
object_types = self.fetch_object_types(cur)
taxonomy_rows = [self.classify_object_type(name) for name in object_types]
capability_rows = self.build_capabilities(taxonomy_rows)
self.create_taxonomy_table(cur, taxonomy_rows)
self.create_capabilities_table(cur, capability_rows)
self.create_detection_view(cur)
self.create_capability_views(cur)
self.ensure_query_indexes(cur)
conn.commit()
print(
"Built navsea_object_taxonomy, navsea_object_capabilities, "
"navsea_detection_objects, capability views and supporting query indexes"
)
@staticmethod
def fetch_object_types(cur) -> list[str]:
cur.execute(
"""
SELECT canonical_object_type
FROM canonical_object_rules
ORDER BY canonical_object_type
"""
)
return [row[0] for row in cur.fetchall()]
@staticmethod
def is_depth_band(name: str) -> bool:
compact = name.replace(" ", "")
return compact.endswith("m") or "m以深" in compact
@staticmethod
def classify_object_type(name: str) -> tuple[str, str, str, str]:
if name.startswith(("P", "L", "p")) and (
name.endswith("ククリ")
or name in {
"P穴",
"P陸域",
"P危険界ククリ",
"P投錨注意障害物ククリ",
"P施設・境界線等ククリ",
"P航路ククリ",
"P錨泊地等ククリ",
"P754ククリ",
"P721ククリ",
"P730ククリ",
"L海底地形",
"p地名",
"p地名陸",
"p高さ制限",
}
):
return (
name,
"overlay_support",
"render_support",
"保留原始 source_layer 的叠加支持对象,不作为独立语义实体替换原始渲染层。",
)
if DetectionPipeline.is_depth_band(name):
return (name, "bathymetry", "depth_zone", "水深分带对象,用于浅滩、深度范围和搁浅风险分析。")
if any(token in name for token in ("等深線",)):
return (name, "bathymetry", "depth_contour", "等深线对象,用于水深变化与航线安全分析。")
if any(token in name for token in ("底質", "海底地形", "サンドウェーブ", "海底火山")):
return (name, "seabed", "seabed_feature", "海底形态或底质对象,用于海底环境识别。")
if any(
token in name
for token in (
"険悪物",
"障害物",
"危険物",
"危険全沈没船",
"沈船",
"全沈没船",
"沈木",
"暗岩",
"洗岩",
"干出岩",
"水上岩",
"孤立危険物",
"魚礁",
"浅所危険界",
"撤去跡",
"掃海済み",
"サンゴ礁",
)
):
return (name, "hazard", "obstacle", "对航行存在碰撞或搁浅风险的危险物对象。")
if any(
token in name
for token in (
"",
"灯台",
"灯標",
"灯浮標",
"浮標",
"立標",
"導灯",
"指向灯",
"V-AIS",
"管制信号所",
)
):
return (name, "navigation_aid", "navigation_mark", "航标、灯标或导助航对象。")
if any(token in name for token in ("航路", "進路矢印", "分離通航方式", "誘導線", "指導線")):
return (name, "route", "route_reference", "航路、导向线或通航组织对象。")
if any(
token in name
for token in (
"境界",
"危険界",
"境界線",
"制限区域",
"専用用途海域",
"航泊禁止区域",
"錨泊禁止区域",
"航空機進入区域",
)
):
return (name, "boundary", "boundary_control", "边界、限制区或规则控制对象。")
if "錨泊" in name or "錨地" in name:
return (name, "anchorage", "anchorage", "锚地或锚泊控制对象。")
if any(token in name for token in ("", "養殖場", "漁網", "海草")):
return (name, "fishery", "fishery_area", "渔业、养殖或缠绕风险相关对象。")
if any(
token in name
for token in (
"防波堤",
"潜提",
"桟橋",
"海上バース",
"ドルフィン",
"ポンツーン",
"",
"道路",
"ビル",
"タンク",
"煙突",
"",
"やぐら",
"ケーソン",
"土砂捨て場",
"石油開発台",
)
):
return (name, "infrastructure", "structure", "港口、岸线或陆上海工构造物对象。")
if any(token in name for token in ("海底線", "輸送管", "送電線", "架空線", "放水口", "取水口")):
return (name, "utility", "utility_line", "海底线缆、管道或公用设施对象。")
if any(token in name for token in ("河川域", "湖沼域", "陸上水域", "干潮帯", "未測海域", "平水境界")):
return (name, "water_area", "water_area", "水域、潮滩或未测区域对象。")
if any(
token in name
for token in (
"その他 (記念碑等)",
"陸上顕著物",
"地名",
"山頂",
"港湾",
"漁港",
"マリーナ",
"フィッシャリーナ",
"海の駅",
"漁業協同組合",
"海事関係署",
"税関",
"パイロットステーション",
)
):
return (name, "place", "place_label", "地名、显著地物或港区服务设施类参考对象。")
if any(token in name for token in ("廃棄物捨て場", "貯木場")):
return (name, "infrastructure", "structure", "作业、堆置或弃置区域类设施对象。")
if any(token in name for token in ("無線局", "レーダー局")):
return (name, "monitoring", "monitoring_station", "通信、雷达或监测设施对象。")
if any(token in name for token in ("急潮", "波紋", "激潮", "渦流", "水源方向")):
return (name, "hazard", "current_hazard", "流态、水流或潮流风险相关对象。")
return (name, "unknown", "review_required", "尚未细化归类的对象,需要人工复核。")
@staticmethod
def build_capabilities(
taxonomy_rows: list[tuple[str, str, str, str]]
) -> list[tuple[str, str]]:
capability_rows: set[tuple[str, str]] = set()
for canonical_object_type, object_family, detection_class, _ in taxonomy_rows:
capabilities = set(FAMILY_DEFAULT_CAPABILITIES.get(object_family, {"review_required"}))
if detection_class in {"obstacle"}:
capabilities.add("collision")
if detection_class in {"depth_zone", "depth_contour"}:
capabilities.add("grounding")
if "魚礁" in canonical_object_type or "浅所" in canonical_object_type or "干潮帯" in canonical_object_type:
capabilities.add("grounding")
if any(token in canonical_object_type for token in ("", "養殖場", "漁網", "海草")):
capabilities.add("entangle")
if any(token in canonical_object_type for token in ("防波堤",)):
capabilities.add("wave_barrier")
if any(token in canonical_object_type for token in ("架空線", "送電線", "高さ制限")):
capabilities.add("overhead_clearance")
if any(token in canonical_object_type for token in ("航路", "誘導線", "進路矢印", "分離通航方式")):
capabilities.add("route_reference")
if any(token in canonical_object_type for token in ("錨泊", "錨地")):
capabilities.add("anchorage_reference")
if any(token in canonical_object_type for token in ("境界", "区域", "危険界")):
capabilities.add("boundary_reference")
if any(token in canonical_object_type for token in ("海底線", "輸送管")):
capabilities.add("snag_risk")
if any(token in canonical_object_type for token in ("", "", "浮標", "導灯", "V-AIS", "管制信号所")):
capabilities.add("navigation_mark")
if object_family in {"place", "infrastructure", "monitoring"}:
capabilities.add("landmark_reference")
for capability in sorted(capabilities):
capability_rows.add((canonical_object_type, capability))
return sorted(capability_rows)
@staticmethod
def create_taxonomy_table(cur, rows: list[tuple[str, str, str, str]]) -> None:
cur.execute("DROP TABLE IF EXISTS navsea_object_taxonomy")
cur.execute(
"""
CREATE TABLE navsea_object_taxonomy (
canonical_object_type VARCHAR(191) NOT NULL,
object_family VARCHAR(64) NOT NULL,
detection_class VARCHAR(64) NOT NULL,
description TEXT NULL,
PRIMARY KEY (canonical_object_type),
KEY idx_family_class (object_family, detection_class)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
"""
)
cur.executemany(
"""
INSERT INTO navsea_object_taxonomy (
canonical_object_type,
object_family,
detection_class,
description
)
VALUES (%s, %s, %s, %s)
""",
rows,
)
@staticmethod
def create_capabilities_table(cur, rows: list[tuple[str, str]]) -> None:
cur.execute("DROP TABLE IF EXISTS navsea_object_capabilities")
cur.execute(
"""
CREATE TABLE navsea_object_capabilities (
canonical_object_type VARCHAR(191) NOT NULL,
capability VARCHAR(64) NOT NULL,
PRIMARY KEY (canonical_object_type, capability),
KEY idx_capability (capability, canonical_object_type)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
"""
)
cur.executemany(
"""
INSERT INTO navsea_object_capabilities (
canonical_object_type,
capability
)
VALUES (%s, %s)
""",
rows,
)
@staticmethod
def create_detection_view(cur) -> None:
cur.execute("DROP VIEW IF EXISTS navsea_detection_objects")
cur.execute(
"""
CREATE VIEW navsea_detection_objects AS
SELECT
r.feature_id,
CAST(NULL AS CHAR(1)) AS geometry,
r.z,
r.x,
r.y,
r.geom_type,
r.source_layer,
r.class_name,
r.shape_name,
r.layer_name,
t.canonical_object_type,
t.object_family,
t.detection_class,
c.capability,
r.detection_key
FROM pbf_relayer_candidates r
JOIN navsea_object_taxonomy t
ON r.canonical_object_type = t.canonical_object_type
LEFT JOIN navsea_object_capabilities c
ON r.canonical_object_type = c.canonical_object_type
"""
)
@staticmethod
def create_capability_views(cur) -> None:
views = {
"navsea_collision_objects": "collision",
"navsea_grounding_objects": "grounding",
"navsea_entangle_objects": "entangle",
"navsea_navigation_mark_objects": "navigation_mark",
"navsea_route_reference_objects": "route_reference",
"navsea_boundary_reference_objects": "boundary_reference",
}
for view_name, capability in views.items():
cur.execute(f"DROP VIEW IF EXISTS {view_name}")
cur.execute(
f"""
CREATE VIEW {view_name} AS
SELECT *
FROM navsea_detection_objects
WHERE capability = %s
""",
(capability,),
)
@staticmethod
def add_index_if_missing(cur, table: str, index_name: str, ddl: str) -> None:
cur.execute(f"SHOW INDEX FROM {table} WHERE Key_name=%s", (index_name,))
if cur.fetchone() is None:
cur.execute(ddl)
def ensure_query_indexes(self, cur) -> None:
self.add_index_if_missing(
cur,
"pbf_relayer_candidates",
"idx_navsea_tile",
"ALTER TABLE pbf_relayer_candidates ADD KEY idx_navsea_tile (z, x, y)",
)
self.add_index_if_missing(
cur,
"pbf_relayer_candidates",
"idx_navsea_object",
"ALTER TABLE pbf_relayer_candidates ADD KEY idx_navsea_object (canonical_object_type(100), geom_type)",
)
self.add_index_if_missing(
cur,
"pbf_relayer_candidates",
"idx_navsea_detection_key",
"ALTER TABLE pbf_relayer_candidates ADD KEY idx_navsea_detection_key (detection_key(64))",
)
self.add_index_if_missing(
cur,
"pbf_relayer_candidates",
"idx_navsea_source_layer",
"ALTER TABLE pbf_relayer_candidates ADD KEY idx_navsea_source_layer (source_layer(100), geom_type)",
)
def main() -> None:
DetectionPipeline(DbConfig()).run()
if __name__ == "__main__":
main()