Roman Klimenko
EN

Deduplikering af CluedIn-datasæt i Python

Biblioteket cluedin er et Python SDK til at arbejde med datamanagement-platformen CluedIn. Det giver adgang til CluedIns REST- og GraphQL-API'er.

CluedIn kan importere data fra forskellige kilder, transformere dem og eksportere dem til andre systemer. Med GraphQL-API'et kan vi forespørge data fra Python eller Jupyter notebooks, ligesom vi ville gøre med enhver anden datakilde.

Installer først de nødvendige biblioteker:

%pip install cluedin
%pip install levenshtein
%pip install Metaphone
%pip install pandas

Opret en JSON-fil med loginoplysninger til instansen. Eksemplet antager, at den er tilgængelig på foobar.244.117.198.49.sslip.io:

{
  "domain": "244.117.198.49.sslip.io",
  "org_name": "foobar",
  "user_email": "admin@foobar.com",
  "user_password": "yourStrong(!)Password"
}

Gem filen et privat sted, og læg stien i en miljøvariabel:

export CLUEDIN_CONTEXT=~/.cluedin/cluedin.json

Følgende kode indlæser data fra CluedIn og gemmer dem i en JSON-fil til videre behandling:

import json
import os

import pandas as pd

import cluedin


def flatten_properties(d):
    for k, v in d['properties'].items():
        if k.startswith('property-'):
            d[k[9:]] = v
    del d['properties']
    return d


def get_entities_from_cluedin(entity_type):
    ctx = cluedin.Context.from_json_file(os.environ['CLUEDIN_CONTEXT'])
    ctx.get_token()

    query = """
        query searchEntities($cursor: PagingCursor, $query: String, $pageSize: Int) {
          search(
            query: $query
            sort: FIELDS
            cursor: $cursor
            pageSize: $pageSize
            sortFields: {field: "id", direction: ASCENDING}
          ) {
            totalResults
            cursor
            entries {
              id
              properties(propertyNames:["imdb.name.primaryName"])
            }
          }
        }
    """

    variables = {
        'query': f'entityType:{entity_type}',
        'pageSize': 10_000
    }

    entities = list(
        map(flatten_properties, cluedin.gql.entries(ctx, query, variables)))

    with open('data.json', 'w', encoding='utf-8') as file:
        json.dump(entities, file, ensure_ascii=False,
                  indent=2, sort_keys=False)


if not os.path.exists('data.json'):
    get_entities_from_cluedin('/IMDb/Person')

data = pd.read_json('data.json')
data.set_index('id', inplace=True)
data

Resultatet er en pandas DataFrame med alle entities af typen /IMDb/Person:

data 0

Antallet af iterationer i værste tilfælde

Målet er at finde mulige dubletter blandt værdierne i egenskaben imdb.name.primaryName.

En naiv løsning sammenligner hvert element med alle andre elementer i en indlejret løkke. Det svarer til mængdens kartesiske kvadrat. Et datasæt med en million elementer giver 1.000.000.000.000 par.

Rækkefølgen i hvert par er ligegyldig. Hvis A er en mulig dublet af B, er B også en mulig dublet af A. Derfor kan vi reducere antallet af sammenligninger fra n^2 til n * (n - 1) / 2, altså omtrent en halvering.

Lad os tjekke, om denne beregning er korrekt:

dataset_size = 1_000
temp = range(dataset_size)

# the result of n(n-1) is always even, because either n or n-1 is even
expected_iterations = int(dataset_size * (dataset_size - 1) / 2)
actual_iterations = 0

for i, a in enumerate(temp):
    for b in temp[i + 1:]:
        actual_iterations += 1

difference = expected_iterations - actual_iterations

print(
    f'Expected iterations: {expected_iterations}, actual iterations: {actual_iterations}, difference: {difference}')

Resultatet er:

Expected iterations: 499500, actual iterations: 499500, difference: 0

Eller vi kan bruge funktionen itertools.combinations:

from itertools import combinations_with_replacement

dataset_size = 1_000

df1 = pd.DataFrame({"name": [f"Person {i}" for i in range(dataset_size)]})

new_df = pd.DataFrame(combinations_with_replacement(
    df1.index, 2), columns=['a', 'b'])

new_df[new_df['a'] != new_df['b']]

Resultatet er 499.500 rækker, præcis som forventet.

Eksakt matching

Ved eksakte matches kan vi undgå det enorme antal sammenligninger. Funktionen duplicated i pandas markerer de entities, der har en eksakt dublet:

data['exact_duplicate'] = data.duplicated(
    subset=['imdb.name.primaryName'], keep=False)

data

Koden kører på under 0,1 sekund og markerer alle eksakte dubletter i vores datasæt med 1.000.000 entities:

data 1

Double Metaphone

Med Double Metaphone beregner vi en hashværdi for hvert navn og sammenligner hashværdierne på samme måde som ved eksakte matches:

import metaphone

data['metaphone'] = data['imdb.name.primaryName'].apply(
    lambda x: metaphone.doublemetaphone(x)[0])
data = data[data['metaphone'] != '']

data['metaphone_duplicates_count'] = \
    data[data['metaphone'] != ''] \
        .groupby('metaphone')
        .transform('count')

data.sort_values(['metaphone_duplicates_count', 'metaphone'],
                 ascending=[False, True], inplace=True)

data

Denne beregning tager lidt under ti sekunder på min tre år gamle bærbare computer og giver os følgende resultat:

data 2
data[data['metaphone_duplicates_count'] > 1].nunique()['metaphone']
# 126320

Levenshtein-afstand

For at finde matches med Levenshtein-afstand skal vi beregne afstanden mellem hvert par af navne. En naiv implementering ville tage omkring 285 år på min bærbare computer, så vi er nødt til at reducere antallet af sammenligninger.

Vi reducerer arbejdet ved at udelukke eksakte dubletter, kun sammenligne navnepar, hvor forskellen i længde højst er ét tegn, og bruge flere processer.

Vi normaliserer Levenshtein-afstanden ved at dividere den med længden af den længste af de to strenge. Dermed kan et længere navn afvige med flere tegn og stadig ligge under den samme grænse.

Implementeringen følger også et map-reduce-mønster. Den beregner afstanden for hvert par, gemmer resultaterne for hver entity i en separat fil og samler derefter filerne til én.

# pylint: disable=missing-module-docstring, missing-function-docstring, missing-class-docstring, line-too-long, disable=invalid-name

# deduplication.py

import json
import os
import time
from datetime import timedelta
from hashlib import md5
from multiprocessing import Pool, cpu_count

from Levenshtein import distance

os.makedirs('levenshtein', exist_ok=True)


def find_duplicates(strings):
    a = strings[0]

    md5_hash = md5(a.encode('utf-8')).hexdigest()
    file_path = f'levenshtein/{md5_hash}.json'

    if os.path.exists(file_path):
        return

    length_distance_threshold = 1
    levenshtein_distance_threshold = 0.2

    result = {
        'name': a,
        'duplicates': []
    }

    for b in strings[1:]:
        if abs(len(a) - len(b)) > length_distance_threshold:
            continue

        levenshtein_distance = distance(a, b) / max(len(a), len(b))

        if levenshtein_distance < levenshtein_distance_threshold:
            result['duplicates'].append({
                'name': b,
                'distance': levenshtein_distance
            })

    if len(result['duplicates']) > 0:
        with open(file_path, 'w', encoding='utf-8') as file:
            json.dump(result, file, indent=2)


def main():
    process_start = int(time.time())

    with open('data.json', 'r', encoding='utf-8') as file: # load the dataset
        data = list({x['imdb.name.primaryName'] for x in json.load(file)})

    print(
        f'{timedelta(seconds=time.time() - process_start)} Loaded {len(data)} unique strings from data.json.')

    with Pool(processes=cpu_count()) as pool:
        print(f'{timedelta(seconds=time.time() - process_start)} Starting {cpu_count()} processes to find duplicates...')

        total = len(data)
        thousand_start = time.time()

        for i, _ in enumerate(pool.imap(find_duplicates, [data[i:] for i, _ in enumerate(data)])):
            if i % 1_000 == 0:
                print(
                    f'{timedelta(seconds=time.time() - process_start)} {timedelta(seconds=time.time() - thousand_start)} strings: {i:_} of {total:_} ({i / total * 100:.0f}%)')
                thousand_start = time.time()

        print(f'{timedelta(seconds=time.time() - process_start)} Done!')

    # Reduce

    files = os.listdir('levenshtein')

    duplicates = {}

    for i, file in enumerate(files):
        if i % 1_000 == 0:
            print(
                f'{timedelta(seconds=time.time() - process_start)} {i:_} of {len(files):_} ({i / len(files) * 100:.0f}%)')
        with open(f'levenshtein/{file}', 'r', encoding='utf-8') as f:
            data = json.load(f)

            if data['name'] not in duplicates:
                duplicates[data['name']] = {
                    'count': 0,
                    'duplicates': []
                }

            duplicates[data['name']]['count'] += len(data['duplicates'])
            duplicates[data['name']]['duplicates'].extend(data['duplicates'])

    duplicates = {k: v for k, v in sorted(
        duplicates.items(), key=lambda item: item[1]['count'], reverse=True)}

    for _, v in duplicates.items():
        v['duplicates'] = sorted(
            v['duplicates'], key=lambda item: item['distance'], reverse=False)

    with open('duplicates.json', 'w', encoding='utf-8') as f: # save the result
        json.dump(duplicates, f, indent=2)

    print(len(files))


if __name__ == '__main__':
    main()

Det tager flere timer at behandle en million entities, men det er stadig bedre end 285 år.

Et BK-træ kan reducere antallet af sammenligninger yderligere, men det må blive en anden artikel.

I praksis justerer jeg først betingelserne for matching på et lille udsnit og kører derefter den endelige beregning på hele datasættet.