fix: read map POIs from authenticated FalkorDB

This commit is contained in:
xuelong committed 2026-08-31 01:51:18 -07:00
1 parent 90f75d41ae
commit 712ca1418c
7 files changed
+646 -206

No files matched your search

+191 -185
View File
@@ -3,6 +3,7 @@ from __future__ import annotations
import asyncio
import json
import logging
import math
import re
import time
@@ -14,12 +15,13 @@ from fastapi import APIRouter, Depends, HTTPException
from app.auth import CurrentUser
from app.config import settings
from app.db import get_agent_settings, get_plaza_overview, get_conn
from app.db import get_agent_settings, get_conn, get_plaza_overview
from app.graph_qa_engine import answer_graph_question
from app.llm_client import LlmClient
from app.project_context import ProjectContext, get_project_context
router = APIRouter()
logger = logging.getLogger(__name__)
SPATIAL_GRAPH_NAME = "guiyang_spatial_v1"
# Libo business zone set: 玉屏街道 12 surveyed city zones (A01–A12),
@@ -58,6 +60,35 @@ GRAPH_PLACE_TYPE_CATEGORIES = {
"poi": "其他地点",
}
# Keep bulk FalkorDB map reads deliberately small. Full imported graph nodes
# can contain raw source payloads, image lists and enrichment results; returning
# ``properties(n)`` for thousands of nodes can exhaust the query response or
# time out even though a lightweight graph count succeeds.
GRAPH_MAP_POI_PROJECTION_FIELDS = (
"element_id",
"gaode_poi_id",
"name",
"type_label",
"place_type",
"lng",
"longitude",
"lat",
"latitude",
"address",
"district",
"city",
"province",
"adcode",
"towncode",
"town_name",
"zone_id",
"zone_name",
"map_poi",
"map_source_labels",
"__kg_node_id",
)
GRAPH_MAP_QUERY_TIMEOUT_MS = 30_000
FOOD_ENRICHMENT_KEYS = (
"food_fusion_status",
"food_match_confidence",
@@ -460,6 +491,20 @@ async def weather_hourly(_user: CurrentUser = None):
}
def _map_falkor_client() -> FalkorDB:
"""Use the same configured database credentials for every map read."""
kwargs: dict[str, Any] = {
"host": settings.falkordb_host,
"port": settings.falkordb_port,
"socket_connect_timeout": 5,
# Allow the 30s server-side query timeout to produce a useful error.
"socket_timeout": 40,
}
if settings.falkordb_password:
kwargs["password"] = settings.falkordb_password
return FalkorDB(**kwargs)
def _is_enterprise_travel_graph(graph_name: str) -> bool:
lower = graph_name.lower()
return "baixinghui" in lower or ("travel" in lower and graph_name != SPATIAL_GRAPH_NAME)
@@ -558,10 +603,7 @@ def _build_libo_bus_route_payload(
def _read_libo_bus_routes(graph_name: str) -> dict[str, Any]:
db = FalkorDB(
host=settings.falkordb_host,
port=settings.falkordb_port,
)
db = _map_falkor_client()
try:
graph = db.select_graph(graph_name)
rows = graph.query(
@@ -666,15 +708,15 @@ def _graph_poi_enrichment(
graph_name: str,
element_id: str,
) -> tuple[list[str], dict[str, Any], dict[str, Any]]:
graph = FalkorDB(
host=settings.falkordb_host,
port=settings.falkordb_port,
).select_graph(graph_name)
result = graph.query(
"MATCH (n {element_id:$element_id}) "
"RETURN labels(n), properties(n) LIMIT 1",
{"element_id": element_id},
).result_set
db = _map_falkor_client()
try:
result = db.select_graph(graph_name).query(
"MATCH (n {element_id:$element_id}) "
"RETURN labels(n), properties(n) LIMIT 1",
{"element_id": element_id},
).result_set
finally:
db.close()
if not result:
return [], {}, {}
@@ -699,6 +741,11 @@ def _graph_poi_business_categories(
labels: list[str] | tuple[str, ...] | None,
properties: dict[str, Any],
) -> list[str]:
source_labels = properties.get("map_source_labels")
if properties.get("map_poi") is True and isinstance(source_labels, list):
# Merging semantic/spatial facets must not alter the verified map's
# business-category contract. All labels remain on the graph node.
labels = source_labels
categories = [
GRAPH_POI_LABELS[label]
for label in (labels or [])
@@ -714,10 +761,7 @@ def _graph_poi_business_categories(
def _graph_poi_categories(graph_name: str) -> dict[str, list[str]]:
db = FalkorDB(
host=settings.falkordb_host,
port=settings.falkordb_port,
)
db = _map_falkor_client()
try:
graph = db.select_graph(graph_name)
result: dict[str, list[str]] = {}
@@ -760,21 +804,59 @@ def _graph_json_value(value: Any) -> Any:
def _graph_map_poi_items(rows: list[list[Any]]) -> list[dict[str, Any]]:
items: list[dict[str, Any]] = []
for node_id, labels, raw_properties in rows:
properties = raw_properties if isinstance(raw_properties, dict) else {}
lng = properties.get("lng", properties.get("longitude"))
lat = properties.get("lat", properties.get("latitude"))
for row in rows:
if len(row) == 3:
node_id, labels, raw_properties = row
properties = raw_properties if isinstance(raw_properties, dict) else {}
elif 20 <= len(row) <= len(GRAPH_MAP_POI_PROJECTION_FIELDS) + 2:
node_id, labels, *values = row
fields = GRAPH_MAP_POI_PROJECTION_FIELDS[:len(values)]
properties = dict(zip(fields, values, strict=True))
else:
continue
if properties.get("map_poi") is False:
continue
normalized_labels = [str(label) for label in (labels or [])]
# Normalized BusStop nodes belong to the dedicated bus-route layer.
# The historical PostgreSQL map source never counted those 90 helper
# nodes as business POIs, so excluding them keeps JSON-created projects
# consistent with the existing 2,799-point map contract.
if "BusStop" in normalized_labels and not any(
label in {"FoodPlace", "Hotel", "ScenicSpot", "TransitFacility"}
for label in normalized_labels
):
continue
lng = properties.get("lng")
if lng is None:
lng = properties.get("longitude")
lat = properties.get("lat")
if lat is None:
lat = properties.get("latitude")
if lng is None or lat is None:
continue
categories = _graph_poi_business_categories(labels, properties)
try:
longitude = float(lng)
latitude = float(lat)
except (TypeError, ValueError):
continue
if (
not math.isfinite(longitude)
or not math.isfinite(latitude)
or not -180 <= longitude <= 180
or not -90 <= latitude <= 90
):
continue
categories = _graph_poi_business_categories(normalized_labels, properties)
if not categories:
continue
category = categories[0]
element_id = str(
properties.get("element_id")
or node_id
or properties.get("__kg_node_id")
or properties.get("gaode_poi_id")
or ""
or (node_id if node_id is not None else "")
)
items.append(
{
@@ -784,8 +866,8 @@ def _graph_map_poi_items(rows: list[list[Any]]) -> list[dict[str, Any]]:
"category": category,
"categories": categories or [category],
"place_type": str(properties.get("place_type") or "poi"),
"lng": float(lng),
"lat": float(lat),
"lng": longitude,
"lat": latitude,
"address": str(properties.get("address") or ""),
"district": str(properties.get("district") or ""),
"city": str(properties.get("city") or ""),
@@ -801,25 +883,43 @@ def _graph_map_poi_items(rows: list[list[Any]]) -> list[dict[str, Any]]:
def _read_graph_map_pois(graph_name: str) -> list[dict[str, Any]]:
db = FalkorDB(host=settings.falkordb_host, port=settings.falkordb_port)
db = _map_falkor_client()
try:
graph = db.select_graph(graph_name)
items: list[dict[str, Any]] = []
page_size = 5000
page_size = 2000
offset = 0
while True:
rows = graph.query(
"MATCH (n) "
"WHERE coalesce(n.lng,n.longitude) IS NOT NULL "
"AND coalesce(n.lat,n.latitude) IS NOT NULL "
"RETURN n.__kg_node_id,labels(n),properties(n) "
f"ORDER BY n.name,n.element_id SKIP {offset} LIMIT {page_size}"
"WHERE (n.lng IS NOT NULL OR n.longitude IS NOT NULL) "
"AND (n.lat IS NOT NULL OR n.latitude IS NOT NULL) "
"AND (n.map_poi IS NULL OR n.map_poi = true) "
"RETURN id(n) AS node_id,labels(n),"
"n.element_id,n.gaode_poi_id,n.name,n.type_label,n.place_type,"
"n.lng,n.longitude,n.lat,n.latitude,n.address,n.district,n.city,"
"n.province,n.adcode,n.towncode,n.town_name,n.zone_id,n.zone_name,"
"n.map_poi,n.map_source_labels,n.__kg_node_id "
f"ORDER BY node_id SKIP {offset} LIMIT {page_size}",
timeout=GRAPH_MAP_QUERY_TIMEOUT_MS,
).result_set
items.extend(_graph_map_poi_items(rows))
if len(rows) < page_size:
break
offset += page_size
return items
# Legacy graphs can contain several label facets for the same POI.
# A map point is identified by its stable POI id, not its graph row.
# Preserve all business categories without drawing duplicate points.
unique: dict[str, dict[str, Any]] = {}
for item in items:
previous = unique.get(item["id"])
if previous is None:
unique[item["id"]] = item
else:
previous["categories"] = list(dict.fromkeys(
[*previous["categories"], *item["categories"]]
))
return list(unique.values())
finally:
db.close()
@@ -890,7 +990,7 @@ def _graph_zone_payload(
def _read_graph_map_zones(graph_name: str) -> dict[str, Any]:
db = FalkorDB(host=settings.falkordb_host, port=settings.falkordb_port)
db = _map_falkor_client()
try:
graph = db.select_graph(graph_name)
town_rows = graph.query(
@@ -944,6 +1044,8 @@ def _graph_poi_detail_payload(
[
*_photo_urls(properties.get("photo_urls")),
*_photo_urls(properties.get("image_urls")),
*_photo_urls(properties.get("dianping_shop_image")),
*_photo_urls(properties.get("hotel_image_samples")),
*image_urls,
]
)
@@ -955,8 +1057,12 @@ def _graph_poi_detail_payload(
or ""
)
categories = _graph_poi_business_categories(labels, properties)
lng = properties.get("lng", properties.get("longitude"))
lat = properties.get("lat", properties.get("latitude"))
lng = properties.get("lng")
if lng is None:
lng = properties.get("longitude")
lat = properties.get("lat")
if lat is None:
lat = properties.get("latitude")
return {
"graph_name": graph_name,
"id": element_id,
@@ -1007,7 +1113,7 @@ def _graph_poi_detail_payload(
def _read_graph_poi_detail(graph_name: str, place_id: str) -> dict[str, Any] | None:
db = FalkorDB(host=settings.falkordb_host, port=settings.falkordb_port)
db = _map_falkor_client()
try:
rows = db.select_graph(graph_name).query(
"MATCH (n) WHERE n.element_id=$place_id "
@@ -1033,73 +1139,51 @@ async def map_pois(
context: ProjectContext = Depends(get_project_context),
_user: CurrentUser = None,
):
"""Return lightweight, project-scoped POI points for the knowledge map."""
"""Return map POIs from the current project's FalkorDB graph.
PostgreSQL spatial tables are collection/staging stores. They are not a
runtime fallback because doing so can display stale POIs from a different
import or project version while the graph counters show current data.
"""
graph_name = _resolve_spatial_graph_name(context.graph_name)
s = settings.db_schema
async with get_conn() as conn:
async with conn.cursor() as cur:
await cur.execute(
f"""SELECT p.element_id, p.gaode_poi_id, p.name, p.type_label,
p.place_type, p.lng, p.lat, p.address, p.province,
p.district, p.city, p.adcode, p.towncode, p.town_name,
z.zone_id, z.zone_name
FROM {s}.amap_spatial_pois p
LEFT JOIN {s}.poi_zone_assignments z
ON z.graph_name = p.graph_name
AND z.gaode_poi_id = p.gaode_poi_id
AND z.zone_set_id = %s
WHERE p.graph_name=%s
ORDER BY p.type_label, p.name
LIMIT 100000""",
(LIBO_ZONE_SET_ID, graph_name),
)
rows = await cur.fetchall()
if rows:
graph_categories: dict[str, list[str]] = {}
if graph_name != SPATIAL_GRAPH_NAME:
try:
graph_categories = await asyncio.to_thread(
_graph_poi_categories,
graph_name,
)
except Exception: # noqa: BLE001 - relational points remain usable
graph_categories = {}
items = [
{
"id": row["element_id"],
"gaode_poi_id": row["gaode_poi_id"],
"name": row["name"] or "未命名POI",
"category": row["type_label"] or "其他地点",
"categories": graph_categories.get(
str(row["element_id"]),
[row["type_label"] or "其他地点"],
),
"place_type": row["place_type"] or "poi",
"lng": float(row["lng"]),
"lat": float(row["lat"]),
"address": row["address"] or "",
"province": row["province"] or "",
"district": row["district"] or "",
"city": row["city"] or "",
"adcode": row["adcode"] or "",
"towncode": row["towncode"] or "",
"town_name": row["town_name"] or "",
"zone_id": row["zone_id"] or "",
"zone_name": row["zone_name"] or "",
}
for row in rows
if row.get("lng") is not None and row.get("lat") is not None
]
else:
# A project created entirely from JSON has no project-specific rows in
# the administrative spatial tables. Its map data lives in FalkorDB.
source = "falkordb"
try:
items = await asyncio.to_thread(_read_graph_map_pois, graph_name)
except Exception as exc:
logger.exception(
"map POI read failed: source=falkordb tenant=%s project=%s graph=%s",
context.tenant_id,
context.project_id,
graph_name,
)
raise HTTPException(
status_code=503,
detail="当前项目的地图 POI 读取失败,请检查 FalkorDB 服务和图谱数据",
) from exc
if items:
logger.info(
"map POI read completed: source=%s tenant=%s project=%s graph=%s total=%d",
source,
context.tenant_id,
context.project_id,
graph_name,
len(items),
)
else:
logger.warning(
"map POI read returned no data: source=%s tenant=%s project=%s graph=%s",
source,
context.tenant_id,
context.project_id,
graph_name,
)
category_counts: dict[str, int] = {}
for item in items:
for category in item["categories"]:
category_counts[category] = category_counts.get(category, 0) + 1
return {
"graph_name": graph_name,
"source": source,
"total": len(items),
"categories": [
{"category": category, "count": count}
@@ -1262,95 +1346,17 @@ async def map_poi_detail(
):
"""Return useful detail fields for one POI in the active project."""
graph_name = _resolve_spatial_graph_name(context.graph_name)
s = settings.db_schema
async with get_conn() as conn:
async with conn.cursor() as cur:
await cur.execute(
f"""SELECT 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, photo_urls, source, source_cell_id,
source_resolution, source_scope_adcode,
raw_jsonb,
first_fetched_at, last_fetched_at
FROM {s}.amap_spatial_pois
WHERE graph_name=%s
AND (element_id=%s OR gaode_poi_id=%s)
LIMIT 1""",
(graph_name, place_id, place_id.removeprefix("amap:")),
)
row = await cur.fetchone()
if not row:
graph_detail = await asyncio.to_thread(
_read_graph_poi_detail,
graph_name,
place_id,
try:
detail = await asyncio.to_thread(_read_graph_poi_detail, graph_name, place_id)
except Exception as exc:
logger.exception(
"map POI detail read failed: project=%s graph=%s",
context.project_id, graph_name,
)
if graph_detail is None:
raise HTTPException(status_code=404, detail="当前项目中未找到该POI")
return graph_detail
raw = row.get("raw_jsonb") or {}
if not isinstance(raw, dict):
raw = {}
graph_labels, food_enrichment, hotel_enrichment = await asyncio.to_thread(
_graph_poi_enrichment,
graph_name,
row["element_id"],
)
external_photos = _photo_urls(food_enrichment.get("dianping_shop_image"))
external_photos.extend(_photo_urls(hotel_enrichment.get("hotel_image_samples")))
photo_urls = list(dict.fromkeys([
*_photo_urls(row["photo_urls"]),
*external_photos,
]))
return {
"graph_name": graph_name,
"id": row["element_id"],
"gaode_poi_id": row["gaode_poi_id"],
"name": row["name"] or "未命名POI",
"category": row["type_label"] or "其他地点",
"place_type": row["place_type"] or "poi",
"business_subcategory": raw.get("business_subcategory") or "",
"scenic_type": raw.get("scenic_type") or "",
"scenic_level": raw.get("scenic_level") or "",
"parent_scenic": raw.get("parent_scenic") or "",
"scenic_grade": raw.get("scenic_grade") or "",
"visitor_value": raw.get("visitor_value") or "",
"audit_result": raw.get("audit_result") or "",
"audit_confidence": raw.get("audit_confidence") or "",
"audit_basis": raw.get("audit_basis") or "",
"audit_date": raw.get("audit_date") or "",
"amap_type": row["amap_type"] or "",
"typecode": row["typecode"] or "",
"scan_hit_count": raw.get("scan_hit_count") or 0,
"matched_scan_types": raw.get("matched_scan_types") or [],
"lng": float(row["lng"]),
"lat": float(row["lat"]),
"province": row["province"] or "",
"city": row["city"] or "",
"district": row["district"] or "",
"adcode": row["adcode"] or "",
"business_area": row["business_area"] or "",
"address": row["address"] or "",
"tel": row["tel"] or "",
"open_time": row["open_time"] or "",
"rating": row["rating"] or "",
"cost": row["cost"] or "",
"level": row["level"] or "",
"tags": row["tags"] or "",
"photo_urls": photo_urls,
"source": row["source"] or "",
"source_cell_id": row["source_cell_id"] or "",
"source_resolution": row["source_resolution"],
"source_scope_adcode": row["source_scope_adcode"] or "",
"first_fetched_at": _iso_datetime(row["first_fetched_at"]),
"last_fetched_at": _iso_datetime(row["last_fetched_at"]),
"graph_labels": graph_labels,
"food_enrichment": food_enrichment,
"hotel_enrichment": hotel_enrichment,
}
raise HTTPException(503, "当前项目的 POI 详情读取失败,请检查 FalkorDB 服务") from exc
if detail is None:
raise HTTPException(404, "当前项目中未找到该POI")
return detail
@router.post("/plaza/user-query")