487 lines
17 KiB
Python
487 lines
17 KiB
Python
"""One-time migration of approved Yunyou Libo graph data into PostgreSQL.
|
|
|
|
The graph remains available for visualization. This module creates the
|
|
relational authority records that the generic data center can manage. It is
|
|
idempotent: record UUIDs are derived from graph identities and every write is
|
|
an upsert.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import json
|
|
import uuid
|
|
from collections.abc import Iterable
|
|
from typing import Any
|
|
|
|
import psycopg
|
|
from falkordb import FalkorDB
|
|
from psycopg import sql
|
|
from psycopg.rows import dict_row
|
|
from psycopg.types.json import Jsonb
|
|
|
|
from app.config import settings
|
|
from app.data_platform.schema import project_schema_name
|
|
|
|
PROJECT_ID = "yunyou_libo"
|
|
GRAPH_NAME = "yunyou_libo"
|
|
UUID_NAMESPACE = uuid.UUID("f9a34c5d-9992-4958-adf8-7495f3251a5d")
|
|
|
|
ENTITY_SPECS = (
|
|
("Hotel", "hotel", "酒店"),
|
|
("FoodPlace", "restaurant", "美食"),
|
|
("ScenicSpot", "scenic", "景区"),
|
|
("TransitFacility", "transport", "交通"),
|
|
("BusStop", "bus_stop", "公交站"),
|
|
)
|
|
|
|
|
|
def stable_uuid(*parts: Any) -> uuid.UUID:
|
|
return uuid.uuid5(UUID_NAMESPACE, ":".join(str(part or "") for part in parts))
|
|
|
|
|
|
def text(value: Any) -> str | None:
|
|
if value is None:
|
|
return None
|
|
result = str(value).strip()
|
|
return result or None
|
|
|
|
|
|
def number(value: Any) -> float | None:
|
|
if value in (None, ""):
|
|
return None
|
|
try:
|
|
return float(value)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
|
|
|
|
def integer(value: Any) -> int | None:
|
|
parsed = number(value)
|
|
return int(parsed) if parsed is not None else None
|
|
|
|
|
|
def json_value(value: Any, fallback: Any) -> Any:
|
|
if value in (None, ""):
|
|
return fallback
|
|
if isinstance(value, (dict, list)):
|
|
return value
|
|
try:
|
|
return json.loads(str(value))
|
|
except (TypeError, ValueError, json.JSONDecodeError):
|
|
return fallback
|
|
|
|
|
|
def category_parts(props: dict[str, Any], fallback: str) -> tuple[str, str | None, str | None]:
|
|
raw_parts = [part.strip() for part in str(props.get("amap_type") or "").split(";") if part.strip()]
|
|
return (
|
|
text(props.get("amap_category_l1")) or (raw_parts[0] if raw_parts else fallback),
|
|
text(props.get("amap_category_l2")) or (raw_parts[1] if len(raw_parts) > 1 else None),
|
|
text(props.get("business_subcategory"))
|
|
or text(props.get("amap_category_l3"))
|
|
or (raw_parts[2] if len(raw_parts) > 2 else None),
|
|
)
|
|
|
|
|
|
def graph_identity(props: dict[str, Any]) -> str:
|
|
return str(
|
|
props.get("element_id")
|
|
or props.get("place_id")
|
|
or props.get("gaode_poi_id")
|
|
or props.get("route_id")
|
|
or props.get("name")
|
|
or uuid.uuid4()
|
|
)
|
|
|
|
|
|
def upsert(
|
|
cur: psycopg.Cursor,
|
|
schema_name: str,
|
|
table_name: str,
|
|
values: dict[str, Any],
|
|
) -> None:
|
|
columns = list(values)
|
|
assignments = [
|
|
sql.SQL("{}=EXCLUDED.{}").format(sql.Identifier(column), sql.Identifier(column))
|
|
for column in columns
|
|
if column not in {"id", "created_at"}
|
|
]
|
|
assignments.extend(
|
|
(
|
|
sql.SQL("updated_at=now()"),
|
|
sql.SQL("deleted_at=NULL"),
|
|
sql.SQL("deleted_by=NULL"),
|
|
)
|
|
)
|
|
cur.execute(
|
|
sql.SQL(
|
|
"INSERT INTO {}.{} ({}) VALUES ({}) "
|
|
"ON CONFLICT (id) DO UPDATE SET {}"
|
|
).format(
|
|
sql.Identifier(schema_name),
|
|
sql.Identifier(table_name),
|
|
sql.SQL(", ").join(sql.Identifier(column) for column in columns),
|
|
sql.SQL(", ").join(sql.Placeholder() for _ in columns),
|
|
sql.SQL(", ").join(assignments),
|
|
),
|
|
list(values.values()),
|
|
)
|
|
|
|
|
|
def graph_rows(graph, label: str, batch_size: int = 400) -> Iterable[dict[str, Any]]:
|
|
"""Read a label in bounded pages so large property payloads do not time out."""
|
|
offset = 0
|
|
while True:
|
|
rows = graph.query(
|
|
f"MATCH (n:{label}) RETURN properties(n) SKIP {offset} LIMIT {batch_size}",
|
|
timeout=120_000,
|
|
).result_set
|
|
for row in rows:
|
|
yield dict(row[0])
|
|
if len(rows) < batch_size:
|
|
break
|
|
offset += batch_size
|
|
|
|
|
|
def compact_source_data(props: dict[str, Any], graph_label: str) -> Jsonb:
|
|
return Jsonb(
|
|
{
|
|
"graph_name": GRAPH_NAME,
|
|
"graph_label": graph_label,
|
|
"graph_element_id": text(props.get("element_id")),
|
|
"place_id": text(props.get("place_id")),
|
|
"source": text(props.get("source")),
|
|
"source_name": text(props.get("source_name")),
|
|
"typecode": text(props.get("typecode")),
|
|
"audit_result": text(props.get("audit_result")),
|
|
"audit_confidence": text(props.get("audit_confidence")),
|
|
}
|
|
)
|
|
|
|
|
|
def insert_external_link(
|
|
cur: psycopg.Cursor,
|
|
schema_name: str,
|
|
tenant_id: str,
|
|
entity_id: uuid.UUID,
|
|
platform: str,
|
|
external_id: Any,
|
|
external_name: Any,
|
|
external_url: Any,
|
|
) -> bool:
|
|
identifier = text(external_id)
|
|
if not identifier:
|
|
return False
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"entity_external_links",
|
|
{
|
|
"id": stable_uuid("external", entity_id, platform, identifier),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"platform": platform,
|
|
"external_id": identifier,
|
|
"external_name": text(external_name),
|
|
"external_url": text(external_url),
|
|
},
|
|
)
|
|
return True
|
|
|
|
|
|
def migrate_entities(
|
|
cur: psycopg.Cursor,
|
|
graph,
|
|
schema_name: str,
|
|
tenant_id: str,
|
|
) -> dict[str, int]:
|
|
counts = {
|
|
"poi_entities": 0,
|
|
"entity_external_links": 0,
|
|
"entity_images": 0,
|
|
"hotel_profiles": 0,
|
|
"restaurant_profiles": 0,
|
|
"scenic_profiles": 0,
|
|
"transport_profiles": 0,
|
|
}
|
|
for graph_label, entity_type, category_fallback in ENTITY_SPECS:
|
|
for props in graph_rows(graph, graph_label):
|
|
identity = graph_identity(props)
|
|
entity_id = stable_uuid(PROJECT_ID, entity_type, identity)
|
|
category_l1, category_l2, category_l3 = category_parts(props, category_fallback)
|
|
entity_name = text(props.get("display_name")) or text(props.get("name")) or identity
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"poi_entities",
|
|
{
|
|
"id": entity_id,
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_type": entity_type,
|
|
"name": entity_name,
|
|
"category_l1": category_l1,
|
|
"category_l2": category_l2,
|
|
"category_l3": category_l3,
|
|
"address": text(props.get("address")),
|
|
"district": text(props.get("district")) or "荔波县",
|
|
"adcode": text(props.get("adcode")),
|
|
"phone": text(props.get("tel")),
|
|
"longitude": number(props.get("lng")),
|
|
"latitude": number(props.get("lat")),
|
|
"h3_r9": text(props.get("h3_r9")),
|
|
"h3_r10": text(props.get("h3_r10")),
|
|
"status": "active",
|
|
"version": 1,
|
|
"extra_data": compact_source_data(props, graph_label),
|
|
},
|
|
)
|
|
counts["poi_entities"] += 1
|
|
|
|
if insert_external_link(
|
|
cur,
|
|
schema_name,
|
|
tenant_id,
|
|
entity_id,
|
|
"amap",
|
|
props.get("gaode_poi_id") or props.get("place_id"),
|
|
entity_name,
|
|
props.get("amap_url"),
|
|
):
|
|
counts["entity_external_links"] += 1
|
|
|
|
cover_image = text(props.get("cover_image_url"))
|
|
if cover_image:
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"entity_images",
|
|
{
|
|
"id": stable_uuid("image", entity_id, "cover", cover_image),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"owner_type": "entity",
|
|
"owner_id": None,
|
|
"image_url": cover_image,
|
|
"caption": "封面图",
|
|
"display_order": 0,
|
|
},
|
|
)
|
|
counts["entity_images"] += 1
|
|
|
|
if entity_type == "hotel" and props.get("ctrip_hotel_id"):
|
|
ctrip_id = props.get("ctrip_hotel_id")
|
|
if insert_external_link(
|
|
cur,
|
|
schema_name,
|
|
tenant_id,
|
|
entity_id,
|
|
"ctrip",
|
|
ctrip_id,
|
|
props.get("ctrip_name_cn"),
|
|
props.get("ctrip_url"),
|
|
):
|
|
counts["entity_external_links"] += 1
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"hotel_profiles",
|
|
{
|
|
"id": stable_uuid("hotel-profile", entity_id),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"ctrip_name": text(props.get("ctrip_name_cn")),
|
|
"opened_year": integer(props.get("ctrip_opened_year")),
|
|
"room_count": integer(props.get("ctrip_room_count")),
|
|
"diamond_level": integer(props.get("ctrip_diamond_level")),
|
|
"ctrip_rating": number(props.get("ctrip_rating")),
|
|
"review_count": integer(props.get("ctrip_review_count")),
|
|
"reference_price": number(props.get("room_min_price") or props.get("offer_min_price")),
|
|
"introduction": text(props.get("ctrip_description")),
|
|
},
|
|
)
|
|
counts["hotel_profiles"] += 1
|
|
|
|
if entity_type == "restaurant" and props.get("dianping_shop_id"):
|
|
dianping_id = props.get("dianping_shop_id")
|
|
if insert_external_link(
|
|
cur,
|
|
schema_name,
|
|
tenant_id,
|
|
entity_id,
|
|
"dianping",
|
|
dianping_id,
|
|
props.get("dianping_name"),
|
|
props.get("dianping_url"),
|
|
):
|
|
counts["entity_external_links"] += 1
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"restaurant_profiles",
|
|
{
|
|
"id": stable_uuid("restaurant-profile", entity_id),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"dianping_name": text(props.get("dianping_name")),
|
|
"dianping_category": text(props.get("dianping_category")),
|
|
"ranking_text": text(props.get("dianping_ranking")),
|
|
"business_status": text(props.get("dianping_business_status")) or "unknown",
|
|
"business_hours": text(props.get("dianping_business_hours")),
|
|
"rating": number(props.get("dianping_rating")),
|
|
"review_count": integer(props.get("dianping_review_count")),
|
|
"average_price": number(props.get("dianping_avg_price")),
|
|
},
|
|
)
|
|
counts["restaurant_profiles"] += 1
|
|
|
|
if entity_type == "scenic":
|
|
scenic_level = text(props.get("scenic_level")) or text(props.get("scenic_grade"))
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"scenic_profiles",
|
|
{
|
|
"id": stable_uuid("scenic-profile", entity_id),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"scenic_type": text(props.get("scenic_type")) or category_l3,
|
|
"scenic_level": scenic_level,
|
|
"is_national": bool(scenic_level and "国家" in scenic_level),
|
|
"visitor_value_type": text(props.get("visitor_value")),
|
|
"opening_hours": text(props.get("open_time")),
|
|
"ticket_note": text(props.get("cost")),
|
|
"official_intro": None,
|
|
},
|
|
)
|
|
counts["scenic_profiles"] += 1
|
|
|
|
if entity_type in {"transport", "bus_stop"}:
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"transport_profiles",
|
|
{
|
|
"id": stable_uuid("transport-profile", entity_id),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"entity_id": entity_id,
|
|
"transport_type": text(props.get("station_type")) or category_l3 or category_l2 or category_l1,
|
|
"service_hours": text(props.get("open_time")),
|
|
"route_note": text(props.get("category")),
|
|
},
|
|
)
|
|
counts["transport_profiles"] += 1
|
|
return counts
|
|
|
|
|
|
def migrate_bus_routes(
|
|
cur: psycopg.Cursor,
|
|
graph,
|
|
schema_name: str,
|
|
tenant_id: str,
|
|
) -> dict[str, int]:
|
|
counts = {"bus_routes": 0, "bus_route_stops": 0}
|
|
for props in graph_rows(graph, "BusRoute"):
|
|
identity = graph_identity(props)
|
|
route_id = stable_uuid(PROJECT_ID, "bus_route", identity)
|
|
service_hours = "—".join(
|
|
part for part in (text(props.get("first_bus")), text(props.get("last_bus"))) if part
|
|
) or None
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"bus_routes",
|
|
{
|
|
"id": route_id,
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"route_name": text(props.get("line_name")) or text(props.get("name")) or identity,
|
|
"direction_name": text(props.get("direction")),
|
|
"start_stop_name": text(props.get("start_stop")),
|
|
"end_stop_name": text(props.get("end_stop")),
|
|
"service_hours": service_hours,
|
|
"route_color": None,
|
|
"geometry": Jsonb([]),
|
|
},
|
|
)
|
|
counts["bus_routes"] += 1
|
|
|
|
rows = graph.query(
|
|
"MATCH (r:BusRoute)-[e:STOPS_AT]->(s:BusStop) "
|
|
"RETURN r.element_id, r.route_id, e.sequence, "
|
|
"s.element_id, s.place_id, s.name, s.lng, s.lat",
|
|
timeout=120_000,
|
|
).result_set
|
|
for row in rows:
|
|
graph_route_id = row[0] or row[1]
|
|
graph_stop_id = row[3] or row[4] or row[5]
|
|
route_id = stable_uuid(PROJECT_ID, "bus_route", graph_route_id)
|
|
stop_entity_id = stable_uuid(PROJECT_ID, "bus_stop", graph_stop_id)
|
|
stop_order = integer(row[2]) or 0
|
|
upsert(
|
|
cur,
|
|
schema_name,
|
|
"bus_route_stops",
|
|
{
|
|
"id": stable_uuid("route-stop", route_id, stop_entity_id, stop_order),
|
|
"tenant_id": tenant_id,
|
|
"project_id": PROJECT_ID,
|
|
"route_id": route_id,
|
|
"stop_entity_id": stop_entity_id,
|
|
"stop_name": text(row[5]) or str(graph_stop_id),
|
|
"stop_order": stop_order,
|
|
"longitude": number(row[6]),
|
|
"latitude": number(row[7]),
|
|
},
|
|
)
|
|
counts["bus_route_stops"] += 1
|
|
return counts
|
|
|
|
|
|
def migrate() -> dict[str, int]:
|
|
graph = FalkorDB(
|
|
host=settings.falkordb_host,
|
|
port=settings.falkordb_port,
|
|
password=settings.falkordb_password or None,
|
|
).select_graph(GRAPH_NAME)
|
|
schema_name = project_schema_name(PROJECT_ID)
|
|
with psycopg.connect(settings.database_url, row_factory=dict_row) as conn:
|
|
with conn.cursor() as cur:
|
|
cur.execute(
|
|
sql.SQL(
|
|
"SELECT tenant_id FROM {}.projects WHERE project_id=%s AND status <> 'archived'"
|
|
).format(sql.Identifier(settings.db_schema)),
|
|
(PROJECT_ID,),
|
|
)
|
|
project = cur.fetchone()
|
|
if not project:
|
|
raise RuntimeError(f"Project not found: {PROJECT_ID}")
|
|
tenant_id = str(project["tenant_id"])
|
|
counts = migrate_entities(cur, graph, schema_name, tenant_id)
|
|
counts.update(migrate_bus_routes(cur, graph, schema_name, tenant_id))
|
|
conn.commit()
|
|
return counts
|
|
|
|
|
|
def main() -> None:
|
|
parser = argparse.ArgumentParser(description=__doc__)
|
|
parser.add_argument(
|
|
"--apply",
|
|
action="store_true",
|
|
help="Execute the idempotent migration. Without this flag only help is shown.",
|
|
)
|
|
args = parser.parse_args()
|
|
if not args.apply:
|
|
parser.print_help()
|
|
return
|
|
print(json.dumps(migrate(), ensure_ascii=False, indent=2))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|