fix: bind knowledge maps to explicit regions

This commit is contained in:
xuelong committed 2026-08-26 02:26:24 -07:00
1 parent 329a94d457
commit 90f75d41ae
15 files changed
+1126 -275

No files matched your search

+285
View File
@@ -0,0 +1,285 @@
#!/usr/bin/env python3
"""Export the authoritative province POI store as an importable map project.
The legacy city project keeps its semantic graph in ``guiyang_new2`` and its
map businesses in PostgreSQL ``amap_spatial_pois``. This exporter snapshots
the latter exactly, so an imported project preserves the 80,609-POI business
count and can use the same knowledge-map template without a hidden graph-name
redirect.
"""
from __future__ import annotations
import argparse
import json
import sys
from collections import Counter
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
import psycopg
from psycopg.rows import dict_row
ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
sys.path.insert(0, str(ROOT))
from app.config import settings
from app.project_lifecycle import normalize_provision_payload
from scripts.export_yunyou_libo_full_graph_bundle import (
SCHEMA_VERSION,
build_schema,
dump_json,
exported_counts,
jsonable,
safe_prefix,
sha256_file,
)
SOURCE_GRAPH_NAME = "guiyang_spatial_v1"
DEFAULT_PROJECT_ID = "city_map_export_v3"
DEFAULT_DISPLAY_NAME = "城市图谱"
# ``raw_jsonb`` and ``photo_urls`` duplicate data that has already been
# normalized into the columns below. Keeping them is useful for an immutable
# archive, but makes an 80,609-POI browser import several hundred megabytes.
# The portable profile keeps every POI while retaining all fields used by the
# shared map template, search and detail drawer.
PORTABLE_PROPERTY_KEYS = {
"element_id",
"gaode_poi_id",
"name",
"type_label",
"place_type",
"amap_type",
"typecode",
"lng",
"lat",
"province",
"city",
"district",
"adcode",
"business_area",
"address",
"tel",
"open_time",
"rating",
"cost",
"level",
"tags",
"source",
"towncode",
"town_name",
}
CATEGORY_LABELS = {
"景点": "ScenicSpot",
"美食": "FoodPlace",
"酒店": "Hotel",
"商场": "Mall",
"医疗保健": "MedicalPlace",
"交通设施": "TransitFacility",
"生活服务": "LifeServicePlace",
"科教文化": "EducationPlace",
"政府机构": "GovernmentPlace",
"公共设施": "Facility",
"体育休闲": "RecreationPlace",
"商务住宅": "ResidentialPlace",
"公司企业": "EnterprisePlace",
"金融保险": "FinancePlace",
"汽车服务": "AutoServicePlace",
"汽车维修": "AutoRepairPlace",
"汽车销售": "AutoSalesPlace",
"摩托车服务": "MotorcycleServicePlace",
"地名地址": "NamedPlace",
"道路附属": "RoadFacility",
}
def read_nodes(
source_graph_name: str,
*,
profile: str,
) -> tuple[list[dict[str, Any]], Counter[str]]:
nodes: list[dict[str, Any]] = []
categories: Counter[str] = Counter()
with psycopg.connect(settings.database_url, row_factory=dict_row) as conn:
with conn.cursor(name="city_map_export") as cur:
cur.execute(
f"""SELECT *
FROM {settings.db_schema}.amap_spatial_pois
WHERE graph_name=%s
ORDER BY element_id""",
(source_graph_name,),
)
for row in cur:
element_id = str(row["element_id"])
category = str(row.get("type_label") or "其他地点")
business_label = CATEGORY_LABELS.get(category, "BusinessPlace")
properties = {
str(key): jsonable(value)
for key, value in row.items()
if key != "graph_name"
and (profile == "full" or key in PORTABLE_PROPERTY_KEYS)
and (
profile == "full"
or value is not None
and value != ""
and value != []
and value != {}
)
}
if profile == "full":
properties["source_graph_name"] = source_graph_name
nodes.append(
{
"id": element_id,
"type": business_label,
"labels": ["Place", business_label],
"properties": properties,
}
)
categories[category] += 1
return nodes, categories
def export_city_map(
output_dir: Path,
*,
project_id: str,
display_name: str,
source_graph_name: str = SOURCE_GRAPH_NAME,
profile: str = "full",
) -> dict[str, Any]:
output_dir.mkdir(parents=True, exist_ok=True)
generated_at = datetime.now(timezone.utc).isoformat()
nodes, category_counts = read_nodes(source_graph_name, profile=profile)
if len(nodes) != len({item["id"] for item in nodes}):
raise RuntimeError("省域 POI 数据存在重复 element_id")
relations: list[dict[str, Any]] = []
graph_name = project_id
spatial_map = {
"enabled": True,
"scope": "guizhou",
"region_name": "贵阳市",
"region_adcode": "520100",
"region_level": "city",
}
counts = exported_counts(nodes, relations)
schema = build_schema(
nodes,
relations,
project_id=project_id,
graph_name=graph_name,
display_name=display_name,
)
source_counts = {
"nodes": len(nodes),
"relations": 0,
"coordinate_nodes": sum(
1
for item in nodes
if item["properties"].get("lng") is not None
and item["properties"].get("lat") is not None
),
"category_counts": dict(sorted(category_counts.items())),
}
graph_data = {
"_bundle": {
"format": "znkg-city-map-snapshot-v3",
"source": "postgresql.amap_spatial_pois",
"source_graph_name": source_graph_name,
"profile": profile,
"project_id": project_id,
"graph_name": graph_name,
"generated_at": generated_at,
"spatial_map": spatial_map,
"source_snapshot_counts": source_counts,
},
"nodes": nodes,
"relations": relations,
}
bundle = {
"format": "znkg-project-bundle-v3",
"project_id": project_id,
"display_name": display_name,
"spatial_map": spatial_map,
"schema": schema,
"graph_data": graph_data,
}
normalized = normalize_provision_payload(bundle)
if (
normalized["counts"]["nodes"] != len(nodes)
or normalized["counts"]["relations"] != 0
):
raise RuntimeError("后端校验后的省域 POI 数量不一致")
prefix = safe_prefix(project_id)
suffix = "full" if profile == "full" else "portable"
schema_path = output_dir / f"{prefix}_{suffix}_schema.v3.json"
graph_path = output_dir / f"{prefix}_{suffix}_graph_data.v3.json"
bundle_path = output_dir / f"{prefix}_{suffix}_bundle.v3.json"
manifest_path = output_dir / f"{prefix}_{suffix}_manifest.v3.json"
dump_json(schema_path, schema, pretty=True)
dump_json(graph_path, graph_data)
dump_json(bundle_path, bundle)
files = {
path.name: {"bytes": path.stat().st_size, "sha256": sha256_file(path)}
for path in (schema_path, graph_path, bundle_path)
}
manifest = {
"format": "znkg-city-map-manifest-v3",
"project_id": project_id,
"graph_name": graph_name,
"source_graph_name": source_graph_name,
"profile": profile,
"schema_version": SCHEMA_VERSION,
"generated_at": generated_at,
"spatial_map": spatial_map,
"source_snapshot_counts": source_counts,
"export_counts": counts,
"validation": "passed",
"validation_checks": {
"node_count_matches_postgresql": counts["nodes"] == source_counts["nodes"],
"coordinate_count_matches_postgresql": (
counts["coordinate_nodes"] == source_counts["coordinate_nodes"]
),
"node_ids_are_unique": True,
"backend_bundle_validation_passed": True,
},
"files": files,
}
dump_json(manifest_path, manifest, pretty=True)
manifest["manifest_file"] = str(manifest_path)
return manifest
def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--output-dir", type=Path, required=True)
parser.add_argument("--project-id", default=DEFAULT_PROJECT_ID)
parser.add_argument("--display-name", default=DEFAULT_DISPLAY_NAME)
parser.add_argument("--source-graph-name", default=SOURCE_GRAPH_NAME)
parser.add_argument(
"--profile",
choices=("full", "portable"),
default="full",
help=(
"full keeps every source column; portable keeps every POI but "
"omits duplicated raw payloads for browser import."
),
)
args = parser.parse_args()
manifest = export_city_map(
args.output_dir.expanduser().resolve(),
project_id=args.project_id,
display_name=args.display_name,
source_graph_name=args.source_graph_name,
profile=args.profile,
)
print(json.dumps(manifest, ensure_ascii=False, indent=2))
if __name__ == "__main__":
main()
+8
View File
@@ -0,0 +1,8 @@
#!/usr/bin/env python3
"""CLI alias for the generic lossless FalkorDB project snapshot exporter."""
from export_yunyou_libo_full_graph_bundle import main
if __name__ == "__main__":
main()
+165 -61
View File
@@ -1,8 +1,8 @@
#!/usr/bin/env python3
"""Export the current Yunyou Libo FalkorDB graph as a lossless import bundle.
"""Export any live FalkorDB graph as a lossless project import bundle.
The live ``yunyou_libo`` graph is the authority for this export. Earlier
versions rebuilt a different graph from PostgreSQL rows: duplicate AMap POIs
The selected live graph is the authority for this export. Earlier project
fixtures rebuilt a different graph from PostgreSQL rows: duplicate AMap POIs
were split, media and links were modeled twice, and unrelated custom tables
were imported automatically. That made the JSON totals differ from the graph
shown by the project.
@@ -40,8 +40,6 @@ PROJECT_ID = "yunyou_libo"
GRAPH_NAME = "yunyou_libo"
DISPLAY_NAME = "云游荔波"
SCHEMA_VERSION = "3.0.0"
SPATIAL_MAP = {"enabled": True, "scope": "libo"}
LABEL_PRIORITY = (
"Hotel",
"FoodPlace",
@@ -196,46 +194,69 @@ def export_live_graph(graph: Any) -> tuple[list[dict[str, Any]], list[dict[str,
exported_id_by_internal: dict[str, str] = {}
exported_ids: set[str] = set()
for row in graph.query("MATCH (n) RETURN n").result_set:
raw_node = row[0]
source_internal_id = internal_id(raw_node)
properties = properties_of(raw_node)
node_id = str(properties.get("__kg_node_id") or f"falkor:{source_internal_id}")
if node_id in exported_ids:
raise RuntimeError(f"图谱存在重复导出节点 ID:{node_id}")
labels = labels_of(raw_node)
exported_ids.add(node_id)
exported_id_by_internal[source_internal_id] = node_id
nodes.append(
{
"id": node_id,
"type": primary_label(labels),
"labels": labels,
"properties": properties,
}
)
page_size = 5000
offset = 0
while True:
rows = graph.query(
"MATCH (n) RETURN n ORDER BY id(n) "
f"SKIP {offset} LIMIT {page_size}"
).result_set
for row in rows:
raw_node = row[0]
source_internal_id = internal_id(raw_node)
properties = properties_of(raw_node)
node_id = str(properties.get("__kg_node_id") or f"falkor:{source_internal_id}")
if node_id in exported_ids:
raise RuntimeError(f"图谱存在重复导出节点 ID:{node_id}")
labels = labels_of(raw_node)
exported_ids.add(node_id)
exported_id_by_internal[source_internal_id] = node_id
nodes.append(
{
"id": node_id,
"type": primary_label(labels),
"labels": labels,
"properties": properties,
}
)
if len(rows) < page_size:
break
offset += page_size
relations: list[dict[str, Any]] = []
rows = graph.query("MATCH (source)-[relation]->(target) RETURN source, relation, target")
for source, edge, target in rows.result_set:
source_id = exported_id_by_internal.get(internal_id(source))
target_id = exported_id_by_internal.get(internal_id(target))
if not source_id or not target_id:
raise RuntimeError("关系引用了未导出的节点")
relations.append(
{
"type": relation_type_of(edge),
"source": source_id,
"target": target_id,
"properties": properties_of(edge),
}
)
offset = 0
while True:
rows = graph.query(
"MATCH (source)-[relation]->(target) "
"RETURN source, relation, target ORDER BY id(relation) "
f"SKIP {offset} LIMIT {page_size}"
).result_set
for source, edge, target in rows:
source_id = exported_id_by_internal.get(internal_id(source))
target_id = exported_id_by_internal.get(internal_id(target))
if not source_id or not target_id:
raise RuntimeError("关系引用了未导出的节点")
relations.append(
{
"type": relation_type_of(edge),
"source": source_id,
"target": target_id,
"properties": properties_of(edge),
}
)
if len(rows) < page_size:
break
offset += page_size
return nodes, relations
def build_schema(
nodes: list[dict[str, Any]],
relations: list[dict[str, Any]],
*,
project_id: str,
graph_name: str,
display_name: str,
) -> dict[str, Any]:
rows_by_label: dict[str, list[dict[str, Any]]] = defaultdict(list)
primary_type_by_id: dict[str, str] = {}
@@ -271,10 +292,10 @@ def build_schema(
}
return {
"namespace": PROJECT_ID,
"namespace": project_id,
"version": SCHEMA_VERSION,
"display_name": "云游荔波当前图谱快照 Schema v3",
"description": "从当前 yunyou_libo FalkorDB 图谱无损导出,保留多标签节点。",
"display_name": f"{display_name}当前图谱快照 Schema v3",
"description": f"从当前 {graph_name} FalkorDB 图谱无损导出,保留多标签节点。",
"entity_types": entity_types,
"relation_types": relation_types,
}
@@ -302,9 +323,14 @@ def exported_counts(
primary_type_counts[item["type"]] += 1
properties = item["properties"]
has_coordinates = properties.get("lng") is not None and properties.get("lat") is not None
if not has_coordinates:
has_coordinates = (
properties.get("longitude") is not None
and properties.get("latitude") is not None
)
if has_coordinates:
coordinate_nodes += 1
if labels & MAP_POI_LABELS:
if labels & MAP_POI_LABELS or properties.get("place_type") or properties.get("type_label"):
map_poi_nodes += 1
if "BusStop" in labels:
bus_stop_nodes += 1
@@ -328,7 +354,8 @@ def source_counts(graph: Any) -> dict[str, Any]:
relation_count = int(graph.query("MATCH ()-[r]->() RETURN count(r)").result_set[0][0])
coordinate_count = int(
graph.query(
"MATCH (n) WHERE n.lng IS NOT NULL AND n.lat IS NOT NULL RETURN count(n)"
"MATCH (n) WHERE coalesce(n.lng,n.longitude) IS NOT NULL "
"AND coalesce(n.lat,n.latitude) IS NOT NULL RETURN count(n)"
).result_set[0][0]
)
label_counts: Counter[str] = Counter()
@@ -405,18 +432,51 @@ def dump_json(path: Path, value: Any, *, pretty: bool = False) -> None:
handle.write("\n")
def export(output_dir: Path) -> dict[str, Any]:
def safe_prefix(value: str) -> str:
normalized = "".join(
char.lower() if char.isalnum() else "_"
for char in value.strip()
)
return "_".join(part for part in normalized.split("_") if part) or "graph"
def export(
output_dir: Path,
*,
project_id: str = PROJECT_ID,
graph_name: str = GRAPH_NAME,
display_name: str = DISPLAY_NAME,
map_enabled: bool = True,
spatial_scope: str = "libo",
region_name: str | None = None,
region_adcode: str | None = None,
region_level: str | None = None,
output_prefix: str | None = None,
) -> dict[str, Any]:
output_dir.mkdir(parents=True, exist_ok=True)
generated_at = datetime.now(timezone.utc).isoformat()
default_region = (
("荔波县", "522722", "district")
if spatial_scope == "libo"
else ("贵州省", "520000", "province")
)
spatial_map = {
"enabled": map_enabled,
"scope": spatial_scope,
"region_name": (region_name or default_region[0]) if map_enabled else "",
"region_adcode": (region_adcode or default_region[1]) if map_enabled else "",
"region_level": (region_level or default_region[2]) if map_enabled else "province",
}
prefix = safe_prefix(output_prefix or project_id)
client = make_client()
try:
graph_names = {
item.decode("utf-8") if isinstance(item, bytes) else str(item)
for item in client.list_graphs()
}
if GRAPH_NAME not in graph_names:
raise RuntimeError(f"FalkorDB 中不存在图谱 {GRAPH_NAME!r}")
graph = client.select_graph(GRAPH_NAME)
if graph_name not in graph_names:
raise RuntimeError(f"FalkorDB 中不存在图谱 {graph_name!r}")
graph = client.select_graph(graph_name)
live_counts = source_counts(graph)
nodes, relations = export_live_graph(graph)
finally:
@@ -424,15 +484,21 @@ def export(output_dir: Path) -> dict[str, Any]:
counts = exported_counts(nodes, relations)
checks = validate_snapshot(live_counts, counts, nodes, relations)
schema = build_schema(nodes, relations)
schema = build_schema(
nodes,
relations,
project_id=project_id,
graph_name=graph_name,
display_name=display_name,
)
graph_data = {
"_bundle": {
"format": "znkg-falkordb-snapshot-v3",
"source": "falkordb",
"project_id": PROJECT_ID,
"graph_name": GRAPH_NAME,
"project_id": project_id,
"graph_name": graph_name,
"generated_at": generated_at,
"spatial_map": SPATIAL_MAP,
"spatial_map": spatial_map,
"source_snapshot_counts": live_counts,
},
"nodes": nodes,
@@ -440,9 +506,9 @@ def export(output_dir: Path) -> dict[str, Any]:
}
provision = {
"format": "znkg-project-bundle-v3",
"project_id": PROJECT_ID,
"display_name": DISPLAY_NAME,
"spatial_map": SPATIAL_MAP,
"project_id": project_id,
"display_name": display_name,
"spatial_map": spatial_map,
"schema": schema,
"graph_data": graph_data,
}
@@ -452,10 +518,10 @@ def export(output_dir: Path) -> dict[str, Any]:
if normalized["counts"]["relations"] != counts["relations"]:
raise RuntimeError("后端校验后的关系数量不一致")
schema_path = output_dir / "yunyou_libo_full_schema.v3.json"
graph_path = output_dir / "yunyou_libo_full_graph_data.v3.json"
bundle_path = output_dir / "yunyou_libo_full_bundle.v3.json"
manifest_path = output_dir / "yunyou_libo_full_manifest.v3.json"
schema_path = output_dir / f"{prefix}_full_schema.v3.json"
graph_path = output_dir / f"{prefix}_full_graph_data.v3.json"
bundle_path = output_dir / f"{prefix}_full_bundle.v3.json"
manifest_path = output_dir / f"{prefix}_full_manifest.v3.json"
dump_json(schema_path, schema, pretty=True)
dump_json(graph_path, graph_data)
dump_json(bundle_path, provision)
@@ -466,12 +532,12 @@ def export(output_dir: Path) -> dict[str, Any]:
}
manifest = {
"format": "znkg-full-graph-manifest-v3",
"project_id": PROJECT_ID,
"graph_name": GRAPH_NAME,
"project_id": project_id,
"graph_name": graph_name,
"schema_version": SCHEMA_VERSION,
"generated_at": generated_at,
"source": "falkordb",
"spatial_map": SPATIAL_MAP,
"spatial_map": spatial_map,
"source_snapshot_counts": live_counts,
"export_counts": counts,
"validation": "passed",
@@ -491,8 +557,46 @@ def main() -> None:
required=True,
help="Directory for schema, graph-data, bundle and manifest JSON files.",
)
parser.add_argument("--project-id", default=PROJECT_ID)
parser.add_argument("--graph-name", default=GRAPH_NAME)
parser.add_argument("--display-name", default=DISPLAY_NAME)
parser.add_argument(
"--map-enabled",
action=argparse.BooleanOptionalAction,
default=True,
help="Whether the exported project should enable the knowledge-map template.",
)
parser.add_argument(
"--map-scope",
choices=("libo", "guizhou"),
default="libo",
help="Boundary scope used by the single knowledge-map template.",
)
parser.add_argument(
"--output-prefix",
default=None,
help="Optional filename prefix; defaults to project-id.",
)
parser.add_argument("--region-name", default=None)
parser.add_argument("--region-adcode", default=None)
parser.add_argument(
"--region-level",
choices=("province", "city", "district"),
default=None,
)
args = parser.parse_args()
manifest = export(args.output_dir.expanduser().resolve())
manifest = export(
args.output_dir.expanduser().resolve(),
project_id=args.project_id,
graph_name=args.graph_name,
display_name=args.display_name,
map_enabled=args.map_enabled,
spatial_scope=args.map_scope,
region_name=args.region_name,
region_adcode=args.region_adcode,
region_level=args.region_level,
output_prefix=args.output_prefix,
)
print(json.dumps(manifest, ensure_ascii=False, indent=2))