Skip to content

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

proposal = infer_from_source_sample(source, vertex_name="party")
cd examples/18-cross-resource-identity
uv run python discover.py

3. Apply it after review

uv run python discover.py --apply

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.

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()