id: splink-entity-resolution-golden-records
namespace: company.team
inputs:
- id: match_threshold
type: FLOAT
displayName: Match probability threshold
description: Two records join one entity when Splink's match probability is
above this. Lower it to 0.9 to watch relatives get chained into one
entity, which the gate blocks.
defaults: 0.99
- id: min_precision
type: FLOAT
displayName: Minimum precision
description: Share of merged pairs that are true matches, measured on the
clerically reviewed sample.
defaults: 0.98
- id: min_recall
type: FLOAT
displayName: Minimum recall
description: Share of true matching pairs that were merged, measured on the
reviewed sample.
defaults: 0.95
- id: max_cluster_size
type: INT
displayName: Largest allowed entity
description: No real customer has more records than there are sources times a
small margin. A larger cluster means unrelated people were chained
together.
defaults: 6
variables:
output_dir: mdm/customers
tasks:
- id: extract_sources
type: io.kestra.plugin.scripts.python.Script
description: >
Demo extract of three systems that each hold their own copy of the same
customers, with the typos, missing values and transposed letters real
systems have. Also writes the clerically reviewed sample, the subset of
records a person has already linked by hand. Replace this task with your
real extracts.
containerImage: python:3.12-slim
outputFiles:
- crm.csv
- billing.csv
- support.csv
- reviewed_sample.csv
script: |
import csv, random
random.seed(11)
FIRST = ["James","Mary","Robert","Patricia","John","Jennifer","Michael","Linda","David","Elizabeth","William","Barbara","Richard","Susan","Joseph","Jessica","Thomas","Sarah","Charles","Karen","Priya","Arjun","Mei","Wei","Olu","Ama","Lars","Ingrid","Diego","Lucia"]
LAST = ["Smith","Johnson","Williams","Brown","Jones","Garcia","Miller","Davis","Rodriguez","Martinez","Patel","Sharma","Chen","Wang","Okafor","Mensah","Larsen","Berg","Silva","Costa","Nguyen","Kim","Ivanov","Novak","Kowalski"]
CITIES = ["Austin","Denver","Seattle","Boston","Chicago","Miami","Portland","Atlanta"]
def typo(s):
if len(s) < 4 or random.random() < 0.5:
return s
i = random.randrange(1, len(s) - 1)
op = random.random()
if op < 0.33:
return s[:i] + s[i + 1:]
if op < 0.66:
return s[:i] + s[i + 1] + s[i] + s[i + 2:]
return s[:i] + random.choice("aeiou") + s[i + 1:]
fields = ["unique_id", "source", "first_name", "surname", "dob", "email", "city", "postcode", "updated_at"]
writers = {s: csv.DictWriter(open(f"{s}.csv", "w", newline=""), fieldnames=fields) for s in ("crm", "billing", "support")}
for w in writers.values():
w.writeheader()
reviewed = csv.writer(open("reviewed_sample.csv", "w", newline=""))
reviewed.writerow(["unique_id", "person_id"])
n = 0
household = None
for pid in range(1, 1201):
# About one in five people shares a household with the previous person: same surname,
# city and postcode, a different first name, birth date and email. These relatives are
# the classic over-merge, and the reason precision has to be measured.
if household and random.random() < 0.2:
l, city, zip_ = household
else:
l, city, zip_ = random.choice(LAST), random.choice(CITIES), f"{random.randint(10000, 99999)}"
household = (l, city, zip_)
f = random.choice(FIRST)
person = {"dob": f"19{random.randint(55, 99)}-{random.randint(1, 12):02d}-{random.randint(1, 28):02d}",
"email": f"{f}.{l}{random.randint(1, 99)}@{random.choice(['gmail.com', 'yahoo.com', 'outlook.com'])}".lower(),
"city": city, "zip": zip_}
for src in ("crm", "billing", "support"):
if src != "crm" and random.random() < 0.35:
continue
uid = f"{src}-{n}"
n += 1
writers[src].writerow({
"unique_id": uid, "source": src,
"first_name": typo(f) if src != "crm" else f,
"surname": typo(l) if src == "support" else l,
"dob": person["dob"] if random.random() > 0.08 else "",
"email": person["email"] if random.random() > 0.25 else "",
"city": person["city"],
"postcode": person["zip"] if random.random() > 0.15 else "",
"updated_at": f"2026-{random.randint(1, 9):02d}-{random.randint(1, 28):02d}",
})
if pid % 4 == 0:
reviewed.writerow([uid, pid])
print(f"extracted {n} records from 3 sources")
- id: resolve_entities
type: io.kestra.plugin.scripts.python.Script
description: >
Train a Splink model without labels (expectation maximisation), score
candidate pairs, cluster records into entities at the threshold, then
measure precision and recall against the reviewed sample and build one
golden record per entity.
containerImage: python:3.12-slim
beforeCommands:
- pip install -q --root-user-action=ignore splink==5.0.0 pandas==2.3.3
inputFiles:
crm.csv: "{{ outputs.extract_sources.outputFiles['crm.csv'] }}"
billing.csv: "{{ outputs.extract_sources.outputFiles['billing.csv'] }}"
support.csv: "{{ outputs.extract_sources.outputFiles['support.csv'] }}"
reviewed_sample.csv: "{{ outputs.extract_sources.outputFiles['reviewed_sample.csv'] }}"
outputFiles:
- golden_records.csv
- crosswalk.csv
env:
MATCH_THRESHOLD: "{{ inputs.match_threshold }}"
script: |
import json, os
from itertools import combinations
import pandas as pd
import splink.comparison_library as cl
from splink import DuckDBAPI, Linker, SettingsCreator, block_on
df = pd.concat([pd.read_csv(f"{s}.csv", dtype=str) for s in ("crm", "billing", "support")], ignore_index=True)
df = df.where(df.notna(), None)
settings = SettingsCreator(
link_type="dedupe_only",
blocking_rules_to_generate_predictions=[
block_on("surname"), block_on("first_name", "dob"), block_on("email"), block_on("postcode"),
],
comparisons=[
cl.JaroWinklerAtThresholds("first_name", [0.92, 0.85]),
cl.JaroWinklerAtThresholds("surname", [0.92, 0.85]),
cl.ExactMatch("dob"), cl.ExactMatch("email"), cl.ExactMatch("postcode"), cl.ExactMatch("city"),
],
)
db = DuckDBAPI()
linker = Linker(db.register(df[["unique_id", "first_name", "surname", "dob", "email", "city", "postcode"]]), settings)
linker.training.estimate_probability_two_random_records_match([block_on("email")], recall=0.7)
linker.training.estimate_u_using_random_sampling(max_pairs=1e6)
linker.training.estimate_parameters_using_expectation_maximisation(block_on("surname"))
linker.training.estimate_parameters_using_expectation_maximisation(block_on("dob"))
scored = linker.inference.predict(threshold_match_probability=0.001)
clusters = linker.clustering.cluster_pairwise_predictions_at_threshold(
scored, threshold_match_probability=float(os.environ["MATCH_THRESHOLD"])).as_pandas_dataframe()
records = df.merge(clusters[["unique_id", "cluster_id"]], on="unique_id")
# Quality against the reviewed sample: pairs of reviewed records only.
reviewed = pd.read_csv("reviewed_sample.csv", dtype=str).merge(records[["unique_id", "cluster_id"]], on="unique_id")
def pairs(groups):
return {tuple(sorted(p)) for g in groups for p in combinations(g, 2)}
predicted = pairs(reviewed.groupby("cluster_id")["unique_id"].apply(list))
truth = pairs(reviewed.groupby("person_id")["unique_id"].apply(list))
hits = len(predicted & truth)
precision = hits / len(predicted) if predicted else 1.0
recall = hits / len(truth) if truth else 1.0
sizes = records.groupby("cluster_id").size()
# Survivorship: newest non-empty value per attribute, CRM wins a tie.
rank = {"crm": 0, "billing": 1, "support": 2}
records["src_rank"] = records["source"].map(rank)
ordered = records.sort_values(["cluster_id", "updated_at", "src_rank"], ascending=[True, False, True])
golden = ordered.groupby("cluster_id").agg(
{c: "first" for c in ["first_name", "surname", "dob", "email", "city", "postcode"]} | {"unique_id": "count"}
).rename(columns={"unique_id": "source_records"}).reset_index().rename(columns={"cluster_id": "entity_id"})
golden.to_csv("golden_records.csv", index=False)
records[["unique_id", "source", "cluster_id"]].rename(columns={"cluster_id": "entity_id"}).to_csv("crosswalk.csv", index=False)
biggest = records[records.cluster_id == sizes.idxmax()]
out = {
"records": len(records), "entities": int(sizes.size), "largest_entity": int(sizes.max()),
"largest_entity_names": ", ".join(sorted(set(biggest.first_name.fillna("") + " " + biggest.surname.fillna(""))))[:200],
"multi_source_entities": int((records.groupby("cluster_id").source.nunique() > 1).sum()),
"reviewed_pairs": len(truth), "precision": round(precision, 4), "recall": round(recall, 4),
}
print(json.dumps(out, indent=2))
print("::" + json.dumps({"outputs": out}) + "::")
- id: log_quality
type: io.kestra.plugin.core.log.Log
message: >-
{{ outputs.resolve_entities.vars.records }} records resolved into {{
outputs.resolve_entities.vars.entities }} entities ({{
outputs.resolve_entities.vars.multi_source_entities }} span several
systems) at threshold {{ inputs.match_threshold }}. Against {{
outputs.resolve_entities.vars.reviewed_pairs }} reviewed pairs: precision
{{ outputs.resolve_entities.vars.precision }}, recall {{
outputs.resolve_entities.vars.recall }}. Largest entity has {{
outputs.resolve_entities.vars.largest_entity }} records ({{
outputs.resolve_entities.vars.largest_entity_names }}).
- id: quality_gate
type: io.kestra.plugin.core.flow.If
description: Publish golden records only when merging is both accurate and
complete enough, and no entity chains unrelated people together.
condition: >-
{{ outputs.resolve_entities.vars.precision >= inputs.min_precision and
outputs.resolve_entities.vars.recall >= inputs.min_recall and
outputs.resolve_entities.vars.largest_entity <= inputs.max_cluster_size }}
then:
- id: publish_golden_records
type: io.kestra.plugin.core.namespace.UploadFiles
description: One golden record per entity, plus the crosswalk from every source
record to its entity.
namespace: "{{ flow.namespace }}"
filesMap:
"{{ vars.output_dir }}/golden_records.csv": "{{ outputs.resolve_entities.outputFiles['golden_records.csv'] }}"
"{{ vars.output_dir }}/crosswalk.csv": "{{ outputs.resolve_entities.outputFiles['crosswalk.csv'] }}"
- id: record_quality
type: io.kestra.plugin.core.kv.Set
key: customer_entity_resolution_quality
kvType: JSON
value: |
{"entities": {{ outputs.resolve_entities.vars.entities }}, "precision": {{ outputs.resolve_entities.vars.precision }},
"recall": {{ outputs.resolve_entities.vars.recall }}, "threshold": {{ inputs.match_threshold }}, "execution_id": "{{ execution.id }}"}
else:
- id: block_publication
type: io.kestra.plugin.core.execution.Fail
description: Over-merging silently fuses customers, which is far worse than
duplicates. The previous golden records stay in place.
errorMessage: >-
Entity resolution below the bar at threshold {{ inputs.match_threshold
}}: precision {{ outputs.resolve_entities.vars.precision }} (min {{
inputs.min_precision }}), recall {{
outputs.resolve_entities.vars.recall }} (min {{ inputs.min_recall }}),
largest entity {{ outputs.resolve_entities.vars.largest_entity }}
records (max {{ inputs.max_cluster_size }}): {{
outputs.resolve_entities.vars.largest_entity_names }}.
triggers:
- id: nightly
type: io.kestra.plugin.core.trigger.Schedule
cron: "0 3 * * *"