The 44 customers with no recoverable name render as "(SIN NOMBRE)", and
ordering the list by name alone floated all of them to the top — "(" sorts
before every letter — so the first two screens of the customer browser were
nothing but placeholders. Small number, worst possible position.
Adds customers.nameMissing, set by the transform and used as the primary sort
key so those records land at the end of the list. Denormalized rather than
computed in the query because the list is paginated in SQL, so the ordering
has to be expressible as a column.
Applied to the dev DB as an ALTER + UPDATE in place (no truncate), so the
existing loaded data and its FKs were left alone.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
366 lines
14 KiB
Python
366 lines
14 KiB
Python
"""
|
|
Migration plan step 3 (customers): build the unified customer master.
|
|
|
|
Reads staged Parquet and loads `customers` + `customer_legacy_refs` in the
|
|
Prisma-managed MySQL. This is the core of the whole project — one customer
|
|
record shared by both business lines — so every later module (policies,
|
|
properties, transactions) resolves its customer FK through the legacy refs
|
|
written here.
|
|
|
|
Rules come from the reconciliation pass (RECONCILIATION.md):
|
|
- Utilities `DATGRAL` (1172) is the customer master; one Customer each.
|
|
- Insurance `DATGRAL` (1070) links to a utilities customer via its
|
|
`num_util` cross-reference. Matches fold into the existing customer (and
|
|
enrich it with the ID-document fields the utilities master lacks);
|
|
non-matches become new insurance-only customers.
|
|
- `COBRO3` is a charge batch, NOT a customer source -> excluded here
|
|
(it is still read as a name-recovery source, see below).
|
|
|
|
`DATGRAL.NOMBRE` is blank on 266 legacy rows (140 utilities, 126 insurance).
|
|
Names for most of them are recovered from secondary tables — see
|
|
`_NAME_SOURCES` — and `customers.nameSource` records which table each
|
|
recovered name came from.
|
|
|
|
Every legacy row folded in gets a `customer_legacy_refs` row
|
|
(sourceSystem, sourceTable=DATGRAL, legacyId) so the merge is auditable and
|
|
the load is idempotent (re-run = truncate + rebuild).
|
|
|
|
Run: ./.venv/bin/python transform_customers.py
|
|
(reads deploy/.env.dev for the dev DATABASE credentials)
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
from decimal import Decimal, InvalidOperation
|
|
from pathlib import Path
|
|
|
|
import pandas as pd
|
|
|
|
from dbenv import connect, env_arg
|
|
|
|
STG = Path(__file__).parent / "output"
|
|
NULL = "∅"
|
|
NOW = datetime.now(timezone.utc).replace(tzinfo=None) # naive UTC for MySQL DATETIME
|
|
|
|
_TRUE = {"1", "-1", "true", "verdadero", "si", "sí", "activo", "yes", "y", "t"}
|
|
_FALSE = {"0", "false", "falso", "no", "inactivo", "n", "f"}
|
|
|
|
|
|
# --------------------------- source normalization --------------------------- #
|
|
def load(source: str, table: str) -> pd.DataFrame:
|
|
df = pd.read_parquet(STG / source / f"{table}.parquet")
|
|
df = df[[c for c in df.columns if c not in ("_legacy_source_table", "_row_num")]].copy()
|
|
for c in df.columns:
|
|
df[c] = df[c].astype("string").str.strip()
|
|
return df
|
|
|
|
|
|
def s(v) -> str | None:
|
|
"""Cell -> clean string or None (empty/sentinel -> None)."""
|
|
if v is None or pd.isna(v):
|
|
return None
|
|
v = str(v).strip()
|
|
return None if v in ("", NULL, "0000-00-00") else v
|
|
|
|
|
|
def norm_id(v) -> str | None:
|
|
v = s(v)
|
|
if v is None:
|
|
return None
|
|
if v.endswith(".0"): # some numeric ids serialize as "521.0"
|
|
v = v[:-2]
|
|
return v if v not in ("0",) else None
|
|
|
|
|
|
def as_bool(v, default=True) -> int:
|
|
v = s(v)
|
|
if v is None:
|
|
return 1 if default else 0
|
|
lv = v.lower()
|
|
if lv in _TRUE:
|
|
return 1
|
|
if lv in _FALSE:
|
|
return 0
|
|
return 1 if default else 0
|
|
|
|
|
|
def as_date(v):
|
|
v = s(v)
|
|
if v is None:
|
|
return None
|
|
dt = pd.to_datetime(v, errors="coerce")
|
|
if pd.isna(dt):
|
|
return None
|
|
return dt.to_pydatetime()
|
|
|
|
|
|
def as_decimal(v):
|
|
v = s(v)
|
|
if v is None:
|
|
return None
|
|
v = v.replace(",", "")
|
|
try:
|
|
return Decimal(v)
|
|
except (InvalidOperation, ValueError):
|
|
return None
|
|
|
|
|
|
# -------------------------------- name recovery ----------------------------- #
|
|
NO_NAME = "(SIN NOMBRE)"
|
|
|
|
# The old PHP importer skipped blank-name DATGRAL rows outright
|
|
# (jorgecuadros-intra-webapp/src/tools/customerAdapter.php:47,81). That also
|
|
# silently dropped those rows' properties, policies and transactions, because
|
|
# every other adapter resolved its customer FK through the customer_mapping
|
|
# table those skipped rows never got into (customerServiceAdapter.php:45), and
|
|
# customerBalanceAdapter.php:52 defaulted the unmapped ones to customer_id 0.
|
|
# Most blank-name rows are real accounts, so recover the name instead of
|
|
# skipping: 176 of the 257 carry a property, policy or transaction.
|
|
#
|
|
# Per side, most trustworthy source first; a later source only fills ids the
|
|
# earlier ones left unresolved. UTILSEG is the office's own hand-maintained
|
|
# name <-> id cross-reference spanning both lines; the rest are billing runs
|
|
# and policy rows that happen to repeat the customer's name.
|
|
# (label, staged source, table, id column, name column)
|
|
_NAME_SOURCES: dict[str, list[tuple[str, str, str, str, str]]] = {
|
|
"utilities": [
|
|
("UTILSEG", "stg_seguros", "utilseg", "util", "nombre"),
|
|
("IVA 2015", "stg_utilities", "iva_2015", "num_id", "nombre"),
|
|
("COBRO3", "stg_utilities", "cobro3", "num_id", "nombre"),
|
|
],
|
|
"insurance": [
|
|
("UTILSEG", "stg_seguros", "utilseg", "seguros", "nombre"),
|
|
("MULT", "stg_seguros", "mult", "num_id", "nombre_aseg"),
|
|
("M EMPR", "stg_seguros", "m_empr", "num_id", "nombre_aseg"),
|
|
("INCENDIO", "stg_seguros", "incendio", "num_id", "nombre_aseg"),
|
|
],
|
|
}
|
|
|
|
|
|
def build_name_index(system: str) -> dict[str, tuple[str, str]]:
|
|
"""legacy num_id -> (recovered name, source label) for one business line."""
|
|
index: dict[str, tuple[str, str]] = {}
|
|
for label, source, table, id_col, name_col in _NAME_SOURCES[system]:
|
|
df = load(source, table)
|
|
if id_col not in df.columns or name_col not in df.columns:
|
|
raise KeyError(f"{table}: expected columns {id_col}/{name_col}, got {list(df.columns)}")
|
|
for _, row in df.iterrows():
|
|
nid, name = norm_id(row[id_col]), s(row[name_col])
|
|
if nid and name and nid not in index:
|
|
index[nid] = (name, label)
|
|
return index
|
|
|
|
|
|
def resolve_name(raw, nid, index) -> tuple[str, str | None]:
|
|
"""(name, nameSource). nameSource stays None when DATGRAL had the name."""
|
|
name = s(raw)
|
|
if name:
|
|
return name, None
|
|
if nid and nid in index:
|
|
return index[nid]
|
|
return NO_NAME, None
|
|
|
|
|
|
# ------------------------------- record builders ---------------------------- #
|
|
def customer_from_utilities(row, name_index) -> dict:
|
|
name, name_source = resolve_name(row["nombre"], norm_id(row["num_id"]), name_index)
|
|
return dict(
|
|
id=str(uuid.uuid4()),
|
|
name=name,
|
|
nameSource=name_source,
|
|
nameMissing=int(name == NO_NAME),
|
|
addressLine1=s(row["direccion"]),
|
|
addressLine2=s(row["colonia"]),
|
|
city=s(row["ciudad"]),
|
|
state=s(row["estado"]),
|
|
zipCode=s(row["codigo"]),
|
|
country=s(row["pais"]),
|
|
phone=s(row["telusa"]),
|
|
mobile=s(row["cel"]),
|
|
fax=s(row["fax"]),
|
|
email=s(row["email"]),
|
|
notes=s(row["observaciones"]),
|
|
identificationType=None,
|
|
identificationNumber=None,
|
|
identificationExpiration=None,
|
|
customerSince=as_date(row["cliente_desde"]),
|
|
status=as_bool(row["status"]),
|
|
feeAmount=as_decimal(row["fee"]),
|
|
updatedAt=NOW,
|
|
)
|
|
|
|
|
|
def customer_from_insurance(row, name_index) -> dict:
|
|
name, name_source = resolve_name(row["nombre"], norm_id(row["num_id"]), name_index)
|
|
return dict(
|
|
id=str(uuid.uuid4()),
|
|
name=name,
|
|
nameSource=name_source,
|
|
nameMissing=int(name == NO_NAME),
|
|
addressLine1=s(row["direccion_1"]),
|
|
addressLine2=s(row["direccion_2"]),
|
|
city=s(row["ciudad"]),
|
|
state=s(row["estado"]),
|
|
zipCode=s(row["codigo"]),
|
|
country=s(row["pais"]),
|
|
phone=s(row["telusa"]),
|
|
mobile=s(row["tel"]),
|
|
fax=s(row["fax"]),
|
|
email=s(row["emailaddress"]),
|
|
notes=s(row["observaciones"]),
|
|
identificationType=s(row["tipo_identificacion"]),
|
|
identificationNumber=s(row["no_identificacion"]),
|
|
identificationExpiration=as_date(row["expira_identificacion"]),
|
|
customerSince=None,
|
|
status=1,
|
|
feeAmount=None,
|
|
updatedAt=NOW,
|
|
)
|
|
|
|
|
|
_CUST_COLS = [
|
|
"id", "name", "nameSource", "nameMissing", "addressLine1", "addressLine2", "city", "state", "zipCode",
|
|
"country", "phone", "mobile", "fax", "email", "notes", "identificationType",
|
|
"identificationNumber", "identificationExpiration", "customerSince",
|
|
"status", "feeAmount", "updatedAt",
|
|
]
|
|
|
|
|
|
def main() -> None:
|
|
env = env_arg()
|
|
conn = connect(env)
|
|
print(f"[customers] target env: {env}")
|
|
cur = conn.cursor()
|
|
|
|
# Fresh, idempotent rebuild.
|
|
cur.execute("SET FOREIGN_KEY_CHECKS=0")
|
|
cur.execute("TRUNCATE TABLE customer_legacy_refs")
|
|
cur.execute("TRUNCATE TABLE customers")
|
|
cur.execute("SET FOREIGN_KEY_CHECKS=1")
|
|
|
|
util = load("stg_utilities", "datgral")
|
|
ins = load("stg_seguros", "datgral")
|
|
util_names = build_name_index("utilities")
|
|
ins_names = build_name_index("insurance")
|
|
|
|
customers: list[dict] = []
|
|
refs: list[tuple] = [] # (id, customerId, sourceSystem, sourceTable, legacyId)
|
|
util_map: dict[str, str] = {} # utilities num_id -> customer_id
|
|
rec_by_id: dict[str, dict] = {}
|
|
|
|
# Phase A: utilities DATGRAL = the master.
|
|
for _, row in util.iterrows():
|
|
rec = customer_from_utilities(row, util_names)
|
|
customers.append(rec)
|
|
rec_by_id[rec["id"]] = rec
|
|
nid = norm_id(row["num_id"])
|
|
legacy = nid or f"rownum_{len(customers)}"
|
|
refs.append((str(uuid.uuid4()), rec["id"], "utilities", "DATGRAL", legacy))
|
|
if nid:
|
|
util_map[nid] = rec["id"]
|
|
|
|
# Phase B: insurance DATGRAL links via num_util, else new customer.
|
|
linked = new_ins = unmatched_numutil = from_ins_side = 0
|
|
enrich: list[tuple] = [] # (customerId, insurance record) for fill-in
|
|
for _, row in ins.iterrows():
|
|
ins_id = norm_id(row["num_id"]) or f"insrow_{new_ins+linked}"
|
|
nutil = norm_id(row["num_util"])
|
|
if nutil and nutil in util_map:
|
|
cust_id = util_map[nutil]
|
|
linked += 1
|
|
ins_rec = customer_from_insurance(row, ins_names)
|
|
# Last name-recovery path: a master whose own line had no name and
|
|
# no utilities-side fallback can still borrow the name its linked
|
|
# insurance record resolved to.
|
|
master = rec_by_id[cust_id]
|
|
if master["name"] == NO_NAME and ins_rec["name"] != NO_NAME:
|
|
master["name"] = ins_rec["name"]
|
|
master["nameSource"] = ins_rec["nameSource"] or "DATGRAL (seguros)"
|
|
master["nameMissing"] = 0
|
|
from_ins_side += 1
|
|
enrich.append((cust_id, ins_rec))
|
|
else:
|
|
if nutil and nutil not in util_map:
|
|
unmatched_numutil += 1
|
|
rec = customer_from_insurance(row, ins_names)
|
|
customers.append(rec)
|
|
rec_by_id[rec["id"]] = rec
|
|
cust_id = rec["id"]
|
|
new_ins += 1
|
|
refs.append((str(uuid.uuid4()), cust_id, "insurance", "DATGRAL", ins_id))
|
|
|
|
# Insert customers.
|
|
placeholders = ",".join(["%s"] * len(_CUST_COLS))
|
|
cur.executemany(
|
|
f"INSERT INTO customers ({','.join(f'`{c}`' for c in _CUST_COLS)}) VALUES ({placeholders})",
|
|
[tuple(rec[c] for c in _CUST_COLS) for rec in customers],
|
|
)
|
|
# Insert legacy refs.
|
|
cur.executemany(
|
|
"INSERT INTO customer_legacy_refs (id, customerId, sourceSystem, sourceTable, legacyId) "
|
|
"VALUES (%s,%s,%s,%s,%s)",
|
|
refs,
|
|
)
|
|
|
|
# Enrich linked customers with insurance-only ID-doc fields, and fill any
|
|
# contact fields the utilities master left empty (COALESCE keeps master's).
|
|
for cust_id, ir in enrich:
|
|
cur.execute(
|
|
"UPDATE customers SET "
|
|
"identificationType = COALESCE(identificationType, %s), "
|
|
"identificationNumber = COALESCE(identificationNumber, %s), "
|
|
"identificationExpiration = COALESCE(identificationExpiration, %s), "
|
|
"email = COALESCE(email, %s), "
|
|
"phone = COALESCE(phone, %s), "
|
|
"mobile = COALESCE(mobile, %s), "
|
|
"updatedAt = %s "
|
|
"WHERE id = %s",
|
|
(ir["identificationType"], ir["identificationNumber"],
|
|
ir["identificationExpiration"], ir["email"], ir["phone"],
|
|
ir["mobile"], NOW, cust_id),
|
|
)
|
|
|
|
conn.commit()
|
|
|
|
# ---- validation / report ----
|
|
cur.execute("SELECT COUNT(*) FROM customers")
|
|
n_cust = cur.fetchone()[0]
|
|
cur.execute("SELECT COUNT(*) FROM customer_legacy_refs")
|
|
n_refs = cur.fetchone()[0]
|
|
cur.execute("SELECT sourceSystem, COUNT(*) FROM customer_legacy_refs GROUP BY sourceSystem")
|
|
by_sys = dict(cur.fetchall())
|
|
cur.execute("SELECT COUNT(*) FROM customer_legacy_refs "
|
|
"GROUP BY customerId HAVING COUNT(*) > 1")
|
|
merged = len(cur.fetchall())
|
|
cur.execute("SELECT nameSource, COUNT(*) FROM customers "
|
|
"WHERE nameSource IS NOT NULL GROUP BY nameSource ORDER BY 2 DESC")
|
|
recovered = cur.fetchall()
|
|
cur.execute("SELECT COUNT(*) FROM customers WHERE name = %s", (NO_NAME,))
|
|
still_unnamed = cur.fetchone()[0]
|
|
|
|
print("=== Customer load complete ===")
|
|
print(f" utilities DATGRAL rows : {len(util)}")
|
|
print(f" insurance DATGRAL rows : {len(ins)}")
|
|
print(f" linked to a utilities customer : {linked}")
|
|
print(f" new insurance-only customers : {new_ins}")
|
|
print(f" num_util set but not in utilities master (data-quality) : {unmatched_numutil}")
|
|
print(f" -> customers : {n_cust} (expected {len(util)} + {new_ins} = {len(util)+new_ins})")
|
|
print(f" -> customer_legacy_refs : {n_refs} (expected {len(util)+len(ins)} = {len(util)+len(ins)})")
|
|
print(f" refs by system : {by_sys}")
|
|
print(f" customers with >1 ref (merged identities) : {merged}")
|
|
print(" name recovery (DATGRAL.NOMBRE was blank):")
|
|
for src, n in recovered:
|
|
print(f" from {src:16} : {n}")
|
|
print(f" of which via the linked insurance record : {from_ins_side}")
|
|
print(f" still {NO_NAME} : {still_unnamed}")
|
|
assert n_cust == len(util) + new_ins, "customer count mismatch"
|
|
assert n_refs == len(util) + len(ins), "legacy ref count mismatch"
|
|
print(" validation: OK")
|
|
conn.close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|