Two systems describe the same customers. How do I find the columns that match them?¶
A CRM export and a billing export describe the same 150 customers. The CRM
calls the email column customer_email, billing calls it email_address; the
CRM says signup_country where billing says country. The values do not
match exactly either: some CRM emails are in upper case, and some billing
emails have spaces around them.
Before you load both exports into one party vertex type, you need to know
which columns hold the same thing, and whether one of them can identify a
customer in both. GraFlo compares every column of one export with every column
of the other, by name and by values, and proposes an identity together with
the evidence for it. Nothing changes until you accept the proposal.
What you need¶
- GraFlo installed (
pip install graflo). - No database is needed, and nothing is written.
The data¶
Both files are generated by generate_data.py with a fixed
seed. Customer names combine first and last names of computer scientists at
random, so "Alan Lovelace" is not a real person.
data/crm_customers.csv, first rows of 150:
| customer_email | full_name | signup_country |
|---|---|---|
| ADA.LOVELACE0@EXAMPLE.COM | Ada Lovelace | GB |
| grace.lovelace1@example.com | Grace Lovelace | US |
| alan.lovelace2@example.com | Alan Lovelace | DE |
data/billing_accounts.csv, first rows of 150:
| email_address | phone | country | invoice_total |
|---|---|---|---|
" ada.lovelace0@example.com " |
+4400892903 | GB | 3461.51 |
| grace.lovelace1@example.com | +73363894879 | US | 567.67 |
| alan.lovelace2@example.com | +20480306579 | DE | 2664.3 |
22 CRM emails are in upper case, and 30 billing emails have spaces around them.
Steps¶
1. Put both exports in one sample¶
discover.py reads the rows of each file into a
ResourceSample, one per resource, and groups them in a SourceSample:
source = SourceSample(
source_name="customer-stack",
description="CRM and billing exports describing the same customers",
samples=[
ResourceSample(resource_name="crm", docs=read_rows("crm_customers.csv")),
ResourceSample(resource_name="billing", docs=read_rows("billing_accounts.csv")),
],
)
For your own sources, GraphEngine.sample_resources builds this sample from
files, a PostgreSQL schema or the bindings of a manifest.
2. Ask for a proposal¶
3. Apply it after review¶
With --apply, the script also passes the proposal to
apply_proposal_to_vertex, which returns a copy of a party vertex with the
proposed identity. The script prints it; nothing is saved.
What you should see¶
discover.py prints:
strategy : natural
identity : ['customer_email']
confidence : 0.83
Column alignments (how the resources were matched up):
left right name values declared
billing.country crm.signup_country 0.67 1.00 False
billing.email_address crm.customer_email 0.37 1.00 False
Per-resource field maps (source -> canonical):
billing: {'country': 'country', 'email_address': 'customer_email'}
crm: {'signup_country': 'country', 'customer_email': 'customer_email'}
Suggested pipeline steps:
{"resource": "billing", "transform": {"rename": {"email_address": "customer_email"}}}
{"resource": "crm", "transform": {"rename": {"signup_country": "country"}}}
Evidence:
doc_counts: {'crm': 150, 'billing': 150}
key_width: 1
resources: ['billing', 'crm']
shared_fields: ['country', 'customer_email']
shared_key_values: 150
uniqueness_by_resource: {'crm': 1.0, 'billing': 1.0}
This is a proposal, not a decision. Nothing was written; review the alignments and evidence before accepting it into a manifest.
Column alignments are the column pairs found to match. name is how alike
the two column names are. values is the overlap of the two value sets after
trimming spaces and lowering case: 1.00 means the same values on both sides. A
pair counts when at least 10% of its values overlap and the average of the two
scores is at least 0.5. The email columns score only 0.37 on names, and pair
because they share every value. declared is False because the pair was found in the
data; a sample taken from PostgreSQL carries its primary and foreign keys, and
those pairs are used as given.
Per-resource field maps give each matched column one name, the
alphabetically first of the pair: customer_email and country.
Suggested pipeline steps are the rename transforms that bring each export
to those names; they go into the pipeline of each resource.
Evidence shows why customer_email is a key: every customer has a
different email within each export (uniqueness_by_resource is 1.0 in both),
and all 150 values appear in both exports (shared_key_values). Uniqueness is
checked in each export separately, because the same customer is meant to
appear in both. country pairs as well, but many customers share a country,
so it is not a key. confidence is the best average score among the
alignments, here the country pair's.
discover.py --apply adds the patched vertex before the closing line:
Patched vertex:
name: party
properties:
- name: full_name
- name: invoice_total
- name: customer_email
identity:
- customer_email
What goes wrong¶
The comparison ignores case and spaces; ingestion does not. In the
proposal, ADA.LOVELACE0@EXAMPLE.COM from the CRM and
" ada.lovelace0@example.com " from billing are the same key. Ingested as they
are, they are two different identities, and Ada becomes two vertices. Before
you ingest with this key, add a transform to both resources that trims and
lowercases the email, for example normalized_key from graflo.util.transform,
as the manifest union example (20) does.
A small sample proves nothing. Each export needs at least 100 rows
(min_sample_size of CrossResourceIdentityConfig); with fewer, the strategy
is no_viable_identity and apply_proposal_to_vertex refuses it.
When no column is shared but each export has its own key, the proposal is an identity funnel with one branch per export, as in the identity funnel example (17). It comes with a warning: each export's records are keyed by that export's own key, so the two records of one customer do not become one vertex.
What to read next¶
- Combine two manifests: two systems that model the same things under different names, matched on a normalized key.
- Cross-resource identity discovery: how a proposal is reached, every strategy and every option.
Files¶
The example lives in examples/18-cross-resource-identity.
discover.py
"""Two systems describe the same customers. How do I find the columns that match them?
Reads the two generated CSV exports in ``data/``, compares their columns by
name and by values, and prints the identity GraFlo proposes for a ``party``
vertex that both describe, with the evidence behind it. With ``--apply`` it
also prints a ``party`` vertex patched with the proposal. Nothing is written.
Run it from this directory:
uv run python discover.py
uv run python discover.py --apply
"""
from __future__ import annotations
import csv
import json
from pathlib import Path
import click
import yaml
from graflo.architecture.onto_sample import ResourceSample, SourceSample
from graflo.architecture.schema.vertex import Vertex
from graflo.db.cross_resource_identity import (
apply_proposal_to_vertex,
infer_from_source_sample,
)
DATA_DIR = Path(__file__).resolve().parent / "data"
def read_rows(name: str) -> list[dict]:
"""Read one CSV file of the example as a list of rows."""
with (DATA_DIR / name).open(encoding="utf-8") as handle:
return list(csv.DictReader(handle))
@click.command()
@click.option(
"--apply",
"apply_",
is_flag=True,
help="Patch a party vertex with the proposal and print the resulting YAML.",
)
def main(apply_: bool) -> None:
"""Sample two CSV resources and propose one identity policy for both."""
source = SourceSample(
source_name="customer-stack",
description="CRM and billing exports describing the same customers",
samples=[
ResourceSample(resource_name="crm", docs=read_rows("crm_customers.csv")),
ResourceSample(
resource_name="billing", docs=read_rows("billing_accounts.csv")
),
],
)
proposal = infer_from_source_sample(source, vertex_name="party")
click.echo(f"strategy : {proposal.strategy}")
click.echo(f"identity : {proposal.identity}")
click.echo(f"confidence : {proposal.confidence:.2f}")
if proposal.warning:
click.echo(f"warning : {proposal.warning}")
click.echo("\nColumn alignments (how the resources were matched up):")
click.echo(
f" {'left':<28} {'right':<28} {'name':>6} {'values':>7} {'declared':>9}"
)
for alignment in proposal.alignments:
left = f"{alignment.left_resource}.{alignment.left_field}"
right = f"{alignment.right_resource}.{alignment.right_field}"
click.echo(
f" {left:<28} {right:<28} "
f"{alignment.name_score:>6.2f} {alignment.value_jaccard:>7.2f} "
f"{alignment.declared!s:>9}"
)
click.echo("\nPer-resource field maps (source -> canonical):")
for resource, mapping in sorted(proposal.resource_field_maps.items()):
click.echo(f" {resource}: {mapping}")
click.echo("\nSuggested pipeline steps:")
for step in proposal.suggested_transforms:
click.echo(f" {json.dumps(step)}")
click.echo("\nEvidence:")
for key, value in sorted(proposal.evidence.items()):
click.echo(f" {key}: {value}")
if apply_:
vertex = Vertex(name="party", properties=["full_name", "invoice_total"])
patched = apply_proposal_to_vertex(vertex, proposal)
click.echo("\nPatched vertex:\n")
click.echo(yaml.safe_dump(patched.to_minimal_canonical_dict(), sort_keys=False))
click.echo(
"\nThis is a proposal, not a decision. Nothing was written; review the "
"alignments and evidence before accepting it into a manifest."
)
if __name__ == "__main__":
main()
generate_data.py
"""
Regenerate the two CSV fixtures for example 18.
Committed output is deterministic (fixed seed), so the discovery result is
stable across runs. Only needed if you want to change the shape of the data:
cd examples/18-cross-resource-identity
uv run python generate_data.py
"""
from __future__ import annotations
import csv
import random
from pathlib import Path
import click
DATA_DIR = Path(__file__).resolve().parent / "data"
N_CUSTOMERS = 150
SEED = 20260731
COUNTRIES = ["GB", "US", "DE", "FR", "NL"]
FIRST = ["ada", "grace", "alan", "edsger", "barbara", "donald", "john", "tony"]
LAST = [
"lovelace",
"hopper",
"turing",
"dijkstra",
"liskov",
"knuth",
"backus",
"hoare",
]
def main() -> None:
rng = random.Random(SEED)
DATA_DIR.mkdir(parents=True, exist_ok=True)
crm_rows = []
billing_rows = []
for i in range(N_CUSTOMERS):
first = FIRST[i % len(FIRST)]
last = LAST[(i // len(FIRST)) % len(LAST)]
email = f"{first}.{last}{i}@example.com"
country = COUNTRIES[i % len(COUNTRIES)]
phone = f"+{rng.randint(1, 99)}{rng.randint(10**8, 10**9 - 1)}"
crm_rows.append(
{
"customer_email": email.upper() if i % 7 == 0 else email,
"full_name": f"{first.title()} {last.title()}",
"signup_country": country,
}
)
# Same people, different column names, plus formatting drift on the
# shared key — which is exactly what `normalize_for_match` absorbs.
billing_rows.append(
{
"email_address": f" {email} " if i % 5 == 0 else email,
"phone": phone,
"country": country,
"invoice_total": round(rng.uniform(10, 5000), 2),
}
)
_write(DATA_DIR / "crm_customers.csv", crm_rows)
_write(DATA_DIR / "billing_accounts.csv", billing_rows)
click.echo(
f"Wrote {N_CUSTOMERS} rows to each of crm_customers.csv, billing_accounts.csv"
)
def _write(path: Path, rows: list[dict]) -> None:
with path.open("w", encoding="utf-8", newline="") as handle:
writer = csv.DictWriter(handle, fieldnames=list(rows[0]))
writer.writeheader()
writer.writerows(rows)
if __name__ == "__main__":
main()