Records arrive with different identifiers filled in. How do I still get one vertex per thing?¶
Your CRM knows customers by email address. Your billing system knows them by
phone number and country, and has an email only for some. You load both into
one vertex type, party. No single column identifies every record: key by
email, and billing rows without an email have no key; key by phone, and CRM
rows have none.
An identity funnel lists the identifiers in order of trust. Each record is keyed by the first identifier it has completely filled in. Two records that share that identifier get the same id, and a database stores them as one vertex, even when they come from different files.
flowchart TD
record([a record]) --> email{email filled in?}
email -- yes --> kemail[id from email]
email -- no --> phone{phone and country filled in?}
phone -- yes --> kphone[id from phone and country]
phone -- no --> weak{name and dob filled in?}
weak -- yes --> kweak[id from name and dob]
weak -- no --> dropped[no id: the record is dropped]
What you need¶
- GraFlo installed (
pip install graflo). - No database is needed. The graph is written to a directory, as in the file backend example (14).
The data¶
| phone | country | name | dob | |
|---|---|---|---|---|
| ada@lovelace.io | Ada Lovelace | 1815-12-10 | ||
| grace@hopper.mil | Grace Hopper | 1906-12-09 | ||
| alan@turing.uk | Alan Turing | 1912-06-23 | ||
| Edsger Dijkstra | 1930-05-11 |
| phone | country | name | dob | |
|---|---|---|---|---|
| +441632960001 | GB | A. Lovelace | 1815-12-10 | |
| +12025550142 | US | G. Hopper | 1906-12-09 | |
| alan@turing.uk | +441632960099 | GB | A. Turing | 1912-06-23 |
| +12025550199 | B. Liskov |
The CRM entered Edsger Dijkstra by hand, without an email. The last billing row has a phone number but no country, and a name but no date of birth: it completes no branch and is meant to be dropped.
Steps¶
1. List the identifiers in order¶
In manifest.yaml, party declares its funnel:
identity: [id]
identity_funnel:
branches:
- id: email
fields: [email]
- id: phone
fields: [phone, country]
- id: weak
fields: [name, dob]
Each branch names the fields that identify a record. GraFlo tries the branches
in order, takes the first whose fields are all filled in, and hashes their
values (SHA-256) into the property id, which is the identity. The branch id
goes into the hash as well, so two branches can never produce the same id.
weak comes last because two different people can share a name and a birth
date.
A single fixed list of fields to hash is declared with
hash_identity_properties; a funnel with one branch does the same. Here one
list is not enough, because no list of fields is filled in on every record.
2. See which identifier keys each record¶
inspect_identities.py runs each row through the
manifest in memory, as ingestion does, and prints the branch that fired and
the id the record received.
3. Write the graph¶
ingest.py loads both files into artifacts/csv-backend, then
counts the party records written and their distinct ids.
What you should see¶
inspect_identities.py prints (it also logs that one record was dropped):
source name branch vertex id
crm Ada Lovelace email 708f5d7d84b9
crm Grace Hopper email ec98fda32ff5
crm Alan Turing email 7f37a4e7c774
crm Edsger Dijkstra weak 9cfd2a53336a
billing A. Lovelace phone b8abd25c5822
billing G. Hopper phone b3dd1f0fc8a0
billing A. Turing email 7f37a4e7c774
billing B. Liskov - none, dropped
8 rows -> 7 records -> 6 distinct ids
Alan Turing has an email in both files, so both of his rows take the email
branch and get the id 7f37a4e7c774.... Edsger Dijkstra falls through to
weak. B. Liskov completes no branch, so the record gets no id and is dropped;
GraFlo does not invent a key, because a random key would add a new vertex on
every run.
ingest.py prints:
Cast dropped 1 'party' document(s) with no value for its identity ['id']. Mark the step lookup_only if the resource only references this vertex.
party: 7 records, 6 distinct ids
The file backend appends records and does not merge them, so Alan's two
records are both in the files, with the same id. A database stores them as
one vertex: six party vertices in all.
What goes wrong¶
The funnel cannot connect different identifiers. Ada Lovelace and Grace
Hopper have an email in the CRM and only a phone number in billing. Their CRM
row takes the email branch and their billing row the phone branch, so
each of them becomes two vertices. Nothing in either record says that
ada@lovelace.io and +441632960001 belong to the same person. Finding that
out needs evidence from the data, such as columns whose values match across
the files; the cross-resource identity example (18)
looks for it.
Changing the funnel changes every id. The branches, their order and their
ids all go into the hash. Change any of them and every record gets a new id,
so the records do not match the vertices already stored. Comparing the two
versions with graflo migrate-schema plan reports this as REKEY_VERTEX with
critical risk, and blocks it.
What to read next¶
- Find matching columns across systems: how GraFlo finds the evidence the funnel lacks.
- Vertex identity:
the funnel,
hash_identity_propertiesand the other identity declarations.
Files¶
The example lives in examples/17-identity-funnel.
manifest.yaml
schema:
metadata:
name: identity-funnel-demo
version: "1.0.0"
graph:
vertex_config:
vertices:
- name: party
properties:
- {name: id, type: STRING}
- {name: email, type: STRING}
- {name: phone, type: STRING}
- {name: country, type: STRING}
- {name: name, type: STRING}
- {name: dob, type: STRING}
identity: [id]
# Branches are tried in order. The first one whose fields are
# all filled in is hashed into id; name and date of birth are
# the weakest evidence, so they come last.
identity_funnel:
branches:
- id: email
fields: [email]
- id: phone
fields: [phone, country]
- id: weak
fields: [name, dob]
edge_config:
edges: []
db_profile: {}
ingestion_model:
resources:
- name: crm
pipeline:
- vertex: party
- name: billing
pipeline:
- vertex: party
bindings:
connectors:
- regex: "^crm\\.csv$"
sub_path: data
resource_name: crm
- regex: "^billing\\.csv$"
sub_path: data
resource_name: billing
ingest.py
"""Records arrive with different identifiers filled in. How do I still get one vertex per thing?
Reads ``manifest.yaml``, loads the CRM and billing files into the file backend
in ``artifacts/csv-backend``, and prints how many ``party`` records were
written and how many distinct ids they carry. Run it from this directory:
uv run python ingest.py
"""
from pathlib import Path
from suthing import FileHandle
from graflo import GraphManifest
from graflo.architecture.backend import GraFloBackendReader
from graflo.connections import GraFloBackendConfig
from graflo.hq import GraphEngine
from graflo.hq.caster import IngestionParams
manifest = GraphManifest.from_config(FileHandle.load("manifest.yaml"))
manifest.finish_init()
backend = GraFloBackendConfig(output_dir=Path("artifacts/csv-backend"))
engine = GraphEngine(target_db_flavor=backend.connection_type)
engine.define_and_ingest(
manifest=manifest,
target_db_config=backend,
ingestion_params=IngestionParams(clear_data=True),
recreate_schema=True,
)
# The file backend appends records; records with the same id are one vertex
# in a database, so count the ids as well.
reader = GraFloBackendReader(backend.output_dir)
parties = [doc for batch in reader.iter_vertex_batches("party") for doc in batch]
distinct = {doc["id"] for doc in parties}
print(f"party: {len(parties)} records, {len(distinct)} distinct ids")
inspect_identities.py
"""Show which identifier keys each customer record, and the vertex id it gets.
Runs every row of ``data/crm.csv`` and ``data/billing.csv`` through
``manifest.yaml`` in memory, as ingestion does, and prints the funnel branch
that fired and the id the record received. No database or file is written.
Run it from this directory:
uv run python inspect_identities.py
"""
import asyncio
import csv
from suthing import FileHandle
from graflo import GraphManifest
from graflo.hq.caster import IngestionParams
from graflo.hq.document_caster import DocumentCaster
ID_WIDTH = 12
manifest = GraphManifest.from_config(FileHandle.load("manifest.yaml"))
manifest.finish_init()
vertices = manifest.require_schema().core_schema.vertex_config.vertices
funnel = next(v for v in vertices if v.name == "party").identity_funnel
assert funnel is not None
caster = DocumentCaster(manifest.require_ingestion_model())
def winning_branch(row: dict[str, str]) -> str:
"""The first branch whose required fields are all filled in, or "-"."""
for branch in funnel.branches:
if all(row.get(field) for field in branch.required_fields):
return branch.id
return "-"
def cast_row(resource: str, row: dict[str, str]) -> dict | None:
"""Cast one row as ingestion does; None when the record is dropped."""
result = asyncio.run(caster.cast_batch([row], resource, params=IngestionParams()))
records = result.graph.vertices.get("party", [])
return records[0] if records else None
rows = records = 0
ids: set[str] = set()
print(f"{'source':<9}{'name':<17}{'branch':<8}vertex id")
for resource in ("crm", "billing"):
with open(f"data/{resource}.csv", newline="", encoding="utf-8") as f:
for row in csv.DictReader(f):
rows += 1
record = cast_row(resource, row)
if record is None:
vertex_id = "none, dropped"
else:
records += 1
ids.add(record["id"])
vertex_id = record["id"][:ID_WIDTH]
print(f"{resource:<9}{row['name']:<17}{winning_branch(row):<8}{vertex_id}")
print(f"{rows} rows -> {records} records -> {len(ids)} distinct ids")