## Linking in Spark


<a target="_blank" href="https://colab.research.google.com/github/moj-analytical-services/splink/blob/splink4_dev/docs/demos/examples/spark/deduplicate_1k_synthetic.ipynb">
  <img src="https://colab.research.google.com/assets/colab-badge.svg" alt="Open In Colab"/>
</a>


In [1]:
# Uncomment and run this cell if you're running in Google Colab.
# !pip install git+https://github.com/moj-analytical-services/splink.git@splink4_dev
# !pip install pyspark

In [2]:
from splink.spark.jar_location import similarity_jar_location

from pyspark import SparkContext, SparkConf
from pyspark.sql import SparkSession

conf = SparkConf()
# This parallelism setting is only suitable for a small toy example
conf.set("spark.driver.memory", "12g")
conf.set("spark.default.parallelism", "16")


# Add custom similarity functions, which are bundled with Splink
# documented here: https://github.com/moj-analytical-services/splink_scalaudfs
path = similarity_jar_location()
conf.set("spark.jars", path)

sc = SparkContext.getOrCreate(conf=conf)

spark = SparkSession(sc)
spark.sparkContext.setCheckpointDir("./tmp_checkpoints")

24/03/13 12:30:29 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable


Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).


In [3]:
# Disable warnings for pyspark - you don't need to include this
import warnings

spark.sparkContext.setLogLevel("ERROR")
warnings.simplefilter("ignore", UserWarning)

In [4]:
from splink import splink_datasets

pandas_df = splink_datasets.fake_1000

df = spark.createDataFrame(pandas_df)

In [5]:
import splink.comparison_library as cl
import splink.comparison_template_library as ctl

from splink import Linker, SparkAPI, block_on, SettingsCreator

settings = SettingsCreator(
    link_type="dedupe_only",
    comparisons=[
        ctl.NameComparison("first_name"),
        ctl.NameComparison("surname"),
        # ctl.date_comparison("dob", cast_strings_to_date=True),
        # TOD=Fix date comparison
        ctl.DateComparison(
            "dob",
            input_is_string=True,
            datetime_metrics=["month", "year", "year"],
            datetime_thresholds=[1, 1, 10],
            datetime_format="%Y%m%d",
        ),
        cl.ExactMatch("city").configure(term_frequency_adjustments=True),
        ctl.EmailComparison("email", include_username_fuzzy_level=False),
    ],
    blocking_rules_to_generate_predictions=[
        block_on("first_name"),
        "l.surname = r.surname",  # alternatively, you can write BRs in their SQL form
    ],
    retain_intermediate_calculation_columns=True,
    em_convergence=0.01,
)

In [6]:
linker = Linker(df, settings, database_api=SparkAPI(spark_session=spark))
deterministic_rules = [
    "l.first_name = r.first_name and levenshtein(r.dob, l.dob) <= 1",
    "l.surname = r.surname and levenshtein(r.dob, l.dob) <= 1",
    "l.first_name = r.first_name and levenshtein(r.surname, l.surname) <= 2",
    "l.email = r.email",
]

linker.estimate_probability_two_random_records_match(deterministic_rules, recall=0.6)

[Stage 0:>                                                        (0 + 12) / 16]



                                                                                



[Stage 3:==>(12 + 4) / 16][Stage 4:>   (0 + 8) / 16][Stage 5:>   (0 + 0) / 16][Stage 4:>  (1 + 12) / 16][Stage 5:>   (0 + 0) / 16][Stage 6:>   (0 + 0) / 16]

[Stage 4:==> (8 + 8) / 16][Stage 5:>   (0 + 4) / 16][Stage 6:>   (0 + 0) / 16]

[Stage 5:>  (0 + 12) / 16][Stage 6:>   (0 + 0) / 16][Stage 8:>    (0 + 0) / 1][Stage 5:>  (4 + 12) / 16][Stage 6:>   (0 + 0) / 16][Stage 8:>    (0 + 0) / 1]

[Stage 5:==>(12 + 4) / 16][Stage 6:>   (0 + 8) / 16][Stage 8:>    (0 + 0) / 1][Stage 6:>  (0 + 12) / 16][Stage 8:>    (0 + 0) / 1][Stage 10:>   (0 + 0) / 1]

[Stage 6:=>  (7 + 9) / 16][Stage 8:>    (0 + 1) / 1][Stage 10:>   (0 + 1) / 1]                                                                                

Probability two random records match is estimated to be  0.0806.
This means that amongst all possible pairwise record comparisons, one in 12.41 are expected to match.  With 499,500 total possible comparisons, we expect a total of around 40,246.67 matching pairs


In [7]:
linker.estimate_u_using_random_sampling(max_pairs=5e5)

----- Estimating u probabilities using random sampling -----










[Stage 50:>                                                       (0 + 12) / 25]







[Stage 57:=> (5 + 8) / 13][Stage 58:>  (0 + 4) / 13][Stage 59:>  (0 + 0) / 13]

[Stage 57:==>(9 + 4) / 13][Stage 58:>  (0 + 8) / 13][Stage 59:>  (0 + 0) / 13]

[Stage 58:> (3 + 10) / 13][Stage 59:>  (0 + 2) / 13][Stage 60:>  (0 + 0) / 13][Stage 58:=> (7 + 6) / 13][Stage 59:>  (0 + 6) / 13][Stage 60:>  (0 + 0) / 13]

[Stage 59:> (2 + 11) / 13][Stage 60:>  (0 + 1) / 13][Stage 62:>  (0 + 0) / 25]

[Stage 60:> (1 + 12) / 13][Stage 62:>  (0 + 0) / 25][Stage 64:>  (0 + 0) / 25]                                                                                

[Stage 72:>               (0 + 12) / 25][Stage 74:>                (0 + 0) / 25]

[Stage 72:> (0 + 12) / 25][Stage 74:>  (0 + 0) / 25][Stage 76:>  (0 + 0) / 25]

[Stage 72:>(13 + 12) / 25][Stage 74:>  (0 + 0) / 25][Stage 76:>  (0 + 0) / 25]                                                                                

[Stage 78:> (0 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:> (3 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:> (5 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:> (7 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:>(10 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:>(12 + 12) / 25][Stage 80:>  (0 + 0) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:>(14 + 11) / 25][Stage 80:>  (0 + 1) / 25][Stage 82:>  (0 + 0) / 25][Stage 78:=>(17 + 8) / 25][Stage 80:>  (3 + 5) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:=>(20 + 5) / 25][Stage 80:=>(16 + 7) / 25][Stage 82:>  (0 + 0) / 25]

[Stage 78:=>(21 + 4) / 25][Stage 82:>  (0 + 8) / 25][Stage 84:>  (0 + 0) / 25][Stage 78:=>(24 + 1) / 25][Stage 82:> (0 + 11) / 25][Stage 84:>  (0 + 0) / 25]

[Stage 78:=>(24 + 1) / 25][Stage 82:> (1 + 11) / 25][Stage 84:>  (0 + 0) / 25][Stage 78:=>(24 + 1) / 25][Stage 82:> (5 + 11) / 25][Stage 84:>  (0 + 0) / 25]

[Stage 78:=>(24 + 1) / 25][Stage 82:>(11 + 11) / 25][Stage 84:>  (0 + 0) / 25]

[Stage 78:=>(24 + 1) / 25][Stage 82:=>(19 + 6) / 25][Stage 84:>  (6 + 5) / 25]



                                                                                


Estimated u probabilities using random sampling



Your model is not yet fully trained. Missing estimates for:
    - first_name (no m values are trained).
    - surname (no m values are trained).
    - dob (no m values are trained).
    - city (no m values are trained).
    - email (no m values are trained).


In [8]:
training_blocking_rule = "l.first_name = r.first_name and l.surname = r.surname"
training_session_fname_sname = (
    linker.estimate_parameters_using_expectation_maximisation(training_blocking_rule)
)

training_blocking_rule = "l.dob = r.dob"
training_session_dob = linker.estimate_parameters_using_expectation_maximisation(
    training_blocking_rule
)


----- Starting EM training session -----



Estimating the m probabilities of the model by blocking on:
l.first_name = r.first_name and l.surname = r.surname

Parameter estimates will be made for the following comparison(s):
    - dob
    - city
    - email

Parameter estimates cannot be made for the following comparison(s) since they are used in the blocking rules: 
    - first_name
    - surname










                                                                                






                                                                                Iteration 1: Largest change in params was -0.697 in probability_two_random_records_match


Iteration 2: Largest change in params was 0.0565 in the m_probability of email, level `All other comparisons`




[Stage 132:>(5 + 10) / 15][Stage 133:> (0 + 2) / 15][Stage 134:> (0 + 0) / 15]

[Stage 133:===>           (3 + 12) / 15][Stage 134:>               (0 + 0) / 15]

                                                                                Iteration 4: Largest change in params was 0.008 in the m_probability of email, level `All other comparisons`



EM converged after 4 iterations



Your model is not yet fully trained. Missing estimates for:
    - first_name (no m values are trained).
    - surname (no m values are trained).



----- Starting EM training session -----



Estimating the m probabilities of the model by blocking on:
l.dob = r.dob

Parameter estimates will be made for the following comparison(s):
    - first_name
    - surname
    - city
    - email

Parameter estimates cannot be made for the following comparison(s) since they are used in the blocking rules: 
    - dob















[Stage 148:>(0 + 12) / 15][Stage 149:> (0 + 0) / 15][Stage 150:> (0 + 0) / 15]

[Stage 148:>(14 + 1) / 15][Stage 149:>(1 + 11) / 15][Stage 150:> (0 + 0) / 15]



Iteration 2: Largest change in params was 0.126 in probability_two_random_records_match




Iteration 4: Largest change in params was 0.0178 in probability_two_random_records_match


Iteration 5: Largest change in params was 0.00899 in probability_two_random_records_match



EM converged after 5 iterations



Your model is fully trained. All comparisons have at least one estimate for their m and u values


In [9]:
results = linker.predict(threshold_match_probability=0.9)

[Stage 203:>                                                      (0 + 12) / 26]

[Stage 203:==>                                                    (1 + 12) / 26][Stage 203:====>                                                  (2 + 12) / 26]



















                                                                                

In [10]:
results.as_pandas_dataframe(limit=5)

Unnamed: 0,match_weight,match_probability,unique_id_l,unique_id_r,first_name_l,first_name_r,gamma_first_name,bf_first_name,surname_l,surname_r,...,gamma_city,tf_city_l,tf_city_r,bf_city,bf_tf_adj_city,email_l,email_r,gamma_email,bf_email,match_key
0,16.4394,0.999989,533,536,Freddie,Freddie,3,11.44555,Gregory,Grerrogy,...,0,0.173,0.187,0.624788,1.0,freddiegregory71@jensen.info,freddiegregory71@jensen.info,3,8.472959,0
1,10.03065,0.999045,838,841,Charlotte,Charlotte,3,11.44555,Harper,Harper,...,0,0.187,0.012,0.624788,1.0,charpe@r@milleb.riz,charper@miller.biz,1,252.763068,0
2,7.476405,0.994416,365,369,Samuel,Samuel,3,11.44555,Campbell,Campbell,...,0,0.001,0.018,0.624788,1.0,samuelcampbell35@hebert.com,samuelcampbell35@hebert.com,3,8.472959,0
3,25.105978,1.0,606,608,Amelia,Amelia,3,11.44555,Porter,Porrret,...,1,0.014,0.014,5.890233,5.089947,ameiiap@nlcholson.orrg,ameliap@nicholson.org,1,252.763068,0
4,8.192765,0.996594,35,416,,,3,11.44555,Bron,Brown,...,0,0.009,0.187,0.624788,1.0,,,3,8.472959,0
