diff --git a/LICENSES/Apache-2.0.txt b/LICENSES/Apache-2.0.txt deleted file mode 100644 index 137069b8..00000000 --- a/LICENSES/Apache-2.0.txt +++ /dev/null @@ -1,73 +0,0 @@ -Apache License -Version 2.0, January 2004 -http://www.apache.org/licenses/ - -TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION - -1. Definitions. - -"License" shall mean the terms and conditions for use, reproduction, and distribution as defined by Sections 1 through 9 of this document. - -"Licensor" shall mean the copyright owner or entity authorized by the copyright owner that is granting the License. - -"Legal Entity" shall mean the union of the acting entity and all other entities that control, are controlled by, or are under common control with that entity. For the purposes of this definition, "control" means (i) the power, direct or indirect, to cause the direction or management of such entity, whether by contract or otherwise, or (ii) ownership of fifty percent (50%) or more of the outstanding shares, or (iii) beneficial ownership of such entity. - -"You" (or "Your") shall mean an individual or Legal Entity exercising permissions granted by this License. - -"Source" form shall mean the preferred form for making modifications, including but not limited to software source code, documentation source, and configuration files. - -"Object" form shall mean any form resulting from mechanical transformation or translation of a Source form, including but not limited to compiled object code, generated documentation, and conversions to other media types. - -"Work" shall mean the work of authorship, whether in Source or Object form, made available under the License, as indicated by a copyright notice that is included in or attached to the work (an example is provided in the Appendix below). - -"Derivative Works" shall mean any work, whether in Source or Object form, that is based on (or derived from) the Work and for which the editorial revisions, annotations, elaborations, or other modifications represent, as a whole, an original work of authorship. For the purposes of this License, Derivative Works shall not include works that remain separable from, or merely link (or bind by name) to the interfaces of, the Work and Derivative Works thereof. - -"Contribution" shall mean any work of authorship, including the original version of the Work and any modifications or additions to that Work or Derivative Works thereof, that is intentionally submitted to Licensor for inclusion in the Work by the copyright owner or by an individual or Legal Entity authorized to submit on behalf of the copyright owner. For the purposes of this definition, "submitted" means any form of electronic, verbal, or written communication sent to the Licensor or its representatives, including but not limited to communication on electronic mailing lists, source code control systems, and issue tracking systems that are managed by, or on behalf of, the Licensor for the purpose of discussing and improving the Work, but excluding communication that is conspicuously marked or otherwise designated in writing by the copyright owner as "Not a Contribution." - -"Contributor" shall mean Licensor and any individual or Legal Entity on behalf of whom a Contribution has been received by Licensor and subsequently incorporated within the Work. - -2. Grant of Copyright License. Subject to the terms and conditions of this License, each Contributor hereby grants to You a perpetual, worldwide, non-exclusive, no-charge, royalty-free, irrevocable copyright license to reproduce, prepare Derivative Works of, publicly display, publicly perform, sublicense, and distribute the Work and such Derivative Works in Source or Object form. - -3. Grant of Patent License. Subject to the terms and conditions of this License, each Contributor hereby grants to You a perpetual, worldwide, non-exclusive, no-charge, royalty-free, irrevocable (except as stated in this section) patent license to make, have made, use, offer to sell, sell, import, and otherwise transfer the Work, where such license applies only to those patent claims licensable by such Contributor that are necessarily infringed by their Contribution(s) alone or by combination of their Contribution(s) with the Work to which such Contribution(s) was submitted. If You institute patent litigation against any entity (including a cross-claim or counterclaim in a lawsuit) alleging that the Work or a Contribution incorporated within the Work constitutes direct or contributory patent infringement, then any patent licenses granted to You under this License for that Work shall terminate as of the date such litigation is filed. - -4. Redistribution. You may reproduce and distribute copies of the Work or Derivative Works thereof in any medium, with or without modifications, and in Source or Object form, provided that You meet the following conditions: - - (a) You must give any other recipients of the Work or Derivative Works a copy of this License; and - - (b) You must cause any modified files to carry prominent notices stating that You changed the files; and - - (c) You must retain, in the Source form of any Derivative Works that You distribute, all copyright, patent, trademark, and attribution notices from the Source form of the Work, excluding those notices that do not pertain to any part of the Derivative Works; and - - (d) If the Work includes a "NOTICE" text file as part of its distribution, then any Derivative Works that You distribute must include a readable copy of the attribution notices contained within such NOTICE file, excluding those notices that do not pertain to any part of the Derivative Works, in at least one of the following places: within a NOTICE text file distributed as part of the Derivative Works; within the Source form or documentation, if provided along with the Derivative Works; or, within a display generated by the Derivative Works, if and wherever such third-party notices normally appear. The contents of the NOTICE file are for informational purposes only and do not modify the License. You may add Your own attribution notices within Derivative Works that You distribute, alongside or as an addendum to the NOTICE text from the Work, provided that such additional attribution notices cannot be construed as modifying the License. - - You may add Your own copyright statement to Your modifications and may provide additional or different license terms and conditions for use, reproduction, or distribution of Your modifications, or for any such Derivative Works as a whole, provided Your use, reproduction, and distribution of the Work otherwise complies with the conditions stated in this License. - -5. Submission of Contributions. Unless You explicitly state otherwise, any Contribution intentionally submitted for inclusion in the Work by You to the Licensor shall be under the terms and conditions of this License, without any additional terms or conditions. Notwithstanding the above, nothing herein shall supersede or modify the terms of any separate license agreement you may have executed with Licensor regarding such Contributions. - -6. Trademarks. This License does not grant permission to use the trade names, trademarks, service marks, or product names of the Licensor, except as required for reasonable and customary use in describing the origin of the Work and reproducing the content of the NOTICE file. - -7. Disclaimer of Warranty. Unless required by applicable law or agreed to in writing, Licensor provides the Work (and each Contributor provides its Contributions) on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied, including, without limitation, any warranties or conditions of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A PARTICULAR PURPOSE. You are solely responsible for determining the appropriateness of using or redistributing the Work and assume any risks associated with Your exercise of permissions under this License. - -8. Limitation of Liability. In no event and under no legal theory, whether in tort (including negligence), contract, or otherwise, unless required by applicable law (such as deliberate and grossly negligent acts) or agreed to in writing, shall any Contributor be liable to You for damages, including any direct, indirect, special, incidental, or consequential damages of any character arising as a result of this License or out of the use or inability to use the Work (including but not limited to damages for loss of goodwill, work stoppage, computer failure or malfunction, or any and all other commercial damages or losses), even if such Contributor has been advised of the possibility of such damages. - -9. Accepting Warranty or Additional Liability. While redistributing the Work or Derivative Works thereof, You may choose to offer, and charge a fee for, acceptance of support, warranty, indemnity, or other liability obligations and/or rights consistent with this License. However, in accepting such obligations, You may act only on Your own behalf and on Your sole responsibility, not on behalf of any other Contributor, and only if You agree to indemnify, defend, and hold each Contributor harmless for any liability incurred by, or claims asserted against, such Contributor by reason of your accepting any such warranty or additional liability. - -END OF TERMS AND CONDITIONS - -APPENDIX: How to apply the Apache License to your work. - -To apply the Apache License to your work, attach the following boilerplate notice, with the fields enclosed by brackets "[]" replaced with your own identifying information. (Don't include the brackets!) The text should be enclosed in the appropriate comment syntax for the file format. We also recommend that a file or class name and description of purpose be included on the same "printed page" as the copyright notice for easier identification within third-party archives. - -Copyright [yyyy] [name of copyright owner] - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - -http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. diff --git a/MANIFEST.in b/MANIFEST.in index 0ea07afa..9f7261c3 100644 --- a/MANIFEST.in +++ b/MANIFEST.in @@ -3,6 +3,5 @@ # SPDX-License-Identifier: MIT include powerplantmatching/package_data/* -include powerplantmatching/package_data/duke_binaries/* include README.md LICENSE include requirements.yaml diff --git a/README.md b/README.md index 843efbe5..fe5ef8ab 100644 --- a/README.md +++ b/README.md @@ -105,7 +105,6 @@ and/or the current release stored on Zenodo with a release-specific DOI: `powerplantmatching` is released as free software under the [MIT](LICENSES/MIT.txt) license. The default output data [powerplants.csv](powerplants.csv) generated by the package is released under [CC BY 4.0](LICENSES/CC-BY-4.0.txt). -Parts of the repository may be licensed under different licenses, especially dependent package binaries for `duke` being licensed under [Apache 2.0 license](https://github.com/PyPSA/powerplantmatching/tree/master/LICENSES/Apache-2.0.txt). This repository uses the [REUSE](https://reuse.software/) conventions to indicate the licenses that apply to individual files and parts of the repository. For details on the licenses that apply, see the the header information of the respective files and [REUSE.toml](REUSE.toml) for details. diff --git a/REUSE.toml b/REUSE.toml index a0f8a6e3..744de935 100644 --- a/REUSE.toml +++ b/REUSE.toml @@ -32,11 +32,3 @@ path = [ ] SPDX-FileCopyrightText = "Contributors to powerplantmatching " SPDX-License-Identifier = "CC0-1.0" - - -[[annotations]] -path = [ - "powerplantmatching/package_data/duke_binaries/*.jar", -] -SPDX-FileCopyrightText = "Lars Marius Garshol " -SPDX-License-Identifier = "Apache-2.0" diff --git a/analysis/benchmark_splink.py b/analysis/benchmark_splink.py new file mode 100644 index 00000000..9662de39 --- /dev/null +++ b/analysis/benchmark_splink.py @@ -0,0 +1,150 @@ +# SPDX-FileCopyrightText: Contributors to powerplantmatching +# +# SPDX-License-Identifier: MIT + +"""Benchmark the Splink matching engine against DUKE on a ground truth +recovered from the production ``powerplants.csv`` (projectID clusters). + +Run from the repository root with the package installed: + + python analysis/benchmark_splink.py +""" + +import ast +import time +import warnings +from itertools import combinations + +import numpy as np +import pandas as pd +from rapidfuzz import fuzz + +warnings.filterwarnings("ignore") +import logging # noqa: E402 + +logging.disable(logging.WARNING) + +import powerplantmatching as pm # noqa: E402 +from powerplantmatching import linkage as lk # noqa: E402 +from powerplantmatching.cleaning import cliques # noqa: E402 + +try: # DUKE was removed in this branch; check out master to benchmark against it + from powerplantmatching.duke import duke +except ImportError: + duke = None + +N_LINK_COUNTRIES = 12 +N_DEDUP_COUNTRIES = 8 + + +def ground_truth(path="powerplants.csv"): + pp = pd.read_csv(path, usecols=["projectID"]) + link, dedup = set(), set() + for raw in pp["projectID"]: + d = ast.literal_eval(raw) + geo, gpd = d.get("GEO", set()), d.get("GPD", set()) + link |= {(g, p) for g in geo for p in gpd} + dedup |= set(combinations(sorted(gpd), 2)) + return link, dedup + + +def prf(pred, truth): + tp = len(pred & truth) + p = tp / len(pred) if pred else 0.0 + r = tp / len(truth) if truth else 0.0 + f = 2 * p * r / (p + r) if p + r else 0.0 + return p, r, f + + +def top_countries(gt, pid2country, available, n): + counts = pd.Series([a for a, _ in gt]).map(pid2country).value_counts() + return [c for c in counts.index if c in available][:n] + + +def link_pairs(matcher, geo, gpd, countries, **kw): + pairs, t0 = set(), time.perf_counter() + for c in countries: + gc, pc = geo[geo.Country == c], gpd[gpd.Country == c] + if gc.empty or pc.empty: + continue + links = matcher([gc, pc], labels=["one", "two"], singlematch=True, **kw) + gp = geo.loc[links["one"], "projectID"].to_numpy() + pp = gpd.loc[links["two"], "projectID"].to_numpy() + pairs |= {(a, b) for a, b in zip(gp, pp)} + return pairs, time.perf_counter() - t0 + + +def _haversine_km(a, b): + la1, lo1, la2, lo2 = map(np.radians, [a.lat, a.lon, b.lat, b.lon]) + h = ( + np.sin((la2 - la1) / 2) ** 2 + + np.cos(la1) * np.cos(la2) * np.sin((lo2 - lo1) / 2) ** 2 + ) + return 2 * 6371 * np.arcsin(np.sqrt(h)) + + +def dedup_pairs(matcher, gpd, countries, **kw): + pairs, merges, principled, total = set(), 0, 0, 0 + t0 = time.perf_counter() + for c in countries: + sub = gpd[gpd.Country == c] + if len(sub) < 2: + continue + groups = cliques(sub, matcher(sub, **kw)) + for _, grp in groups.groupby("grouped"): + ids = sorted(grp.projectID.unique()) + if len(ids) < 2: + continue + merges += 1 + pairs |= set(combinations(ids, 2)) + for a, b in combinations(grp.itertuples(), 2): + total += 1 + name_sim = fuzz.token_set_ratio(str(a.Name), str(b.Name)) / 100 + far = pd.isna(a.lat) or pd.isna(b.lat) or _haversine_km(a, b) > 5 + principled += name_sim >= 0.6 and not far and a.Fueltype == b.Fueltype + return pairs, time.perf_counter() - t0, merges, principled / total if total else 0.0 + + +def main(): + gt_link, gt_dedup = ground_truth() + geo, gpd = pm.data.GEO(), pm.data.GPD() + geo_ids, gpd_ids = set(geo.projectID), set(gpd.projectID) + gt_link = {(g, p) for g, p in gt_link if g in geo_ids and p in gpd_ids} + gt_dedup = {(a, b) for a, b in gt_dedup if a in gpd_ids and b in gpd_ids} + + shared = set(geo.Country) & set(gpd.Country) + geo_c = geo.set_index("projectID").Country.to_dict() + gpd_c = gpd.set_index("projectID").Country.to_dict() + link_countries = top_countries(gt_link, geo_c, shared, N_LINK_COUNTRIES) + dedup_countries = top_countries( + gt_dedup, gpd_c, set(gpd.Country), N_DEDUP_COUNTRIES + ) + + print(f"== record linkage (GEO x GPD, {len(link_countries)} countries) ==") + print(f"ground-truth pairs: {len(gt_link)}") + if duke is not None: + dp, dt = link_pairs(duke, geo, gpd, link_countries) + p, r, f = prf(dp, gt_link) + print(f"DUKE P={p:.3f} R={r:.3f} F1={f:.3f} ({dt:.1f}s)") + for th in [0.5, 0.7, 0.9, 0.95]: + sp, st = link_pairs(lk.match, geo, gpd, link_countries, threshold=th) + p, r, f = prf(sp, gt_link) + print(f"splink {th:.2f} P={p:.3f} R={r:.3f} F1={f:.3f} ({st:.1f}s)") + + print(f"\n== deduplication (GPD, {len(dedup_countries)} countries) ==") + print(f"ground-truth pairs (circular, DUKE-derived): {len(gt_dedup)}") + if duke is not None: + dp, dt, dm, _ = dedup_pairs(duke, gpd, dedup_countries) + p, r, f = prf(dp, gt_dedup) + print(f"DUKE merges={dm} P={p:.3f} R={r:.3f} F1={f:.3f} ({dt:.1f}s)") + for th in [0.9, 0.95, 0.99]: + sp, st, sm, frac = dedup_pairs(lk.match, gpd, dedup_countries, threshold=th) + p, r, f = prf(sp, gt_dedup) + print( + f"splink {th:.2f} merges={sm} P={p:.3f} R={r:.3f} " + f"F1={f:.3f} principled={100 * frac:.0f}% ({st:.1f}s)" + ) + + +if __name__ == "__main__": + main() diff --git a/analysis/linkage_findings.md b/analysis/linkage_findings.md new file mode 100644 index 00000000..0936df4f --- /dev/null +++ b/analysis/linkage_findings.md @@ -0,0 +1,92 @@ + + +# Why the matching engine replaced DUKE (Splink experiment) + +`powerplantmatching` historically matched records with [DUKE](https://github.com/larsga/Duke), +a Java/JVM record-linkage engine invoked as a subprocess and configured via XML. +This branch replaces it with `powerplantmatching.linkage` — a pure-Python +backend built on [Splink](https://moj-analytical-services.github.io/splink) and +DuckDB — which needs no JVM, configures the model in Python, and is several +times faster at comparable quality on an objective ground truth. This note +records the evidence. (A sibling experiment, #301, swapped the same engine for a +`rapidfuzz` + `numpy` backend; this is the Splink variant.) + +## What the engine does + +`linkage.match` takes a list of two frames (record linkage) or a single frame +(deduplication) and returns matched index pairs. It builds a **fixed-parameter +Splink (Fellegi-Sunter) model** — the per-field `m`/`u` probabilities are set +directly from the former DUKE field weights (`Comparison.xml` / +`Deleteduplicates.xml`), so the model needs **no per-call EM training** and is +fully deterministic and stable across the many small per-country slices the +pipeline produces. + +- record linkage ← `Comparison.xml` (Name, Fueltype, Country, Capacity, Geo) +- deduplication ← `Deleteduplicates.xml` (adds Technology, near-neutral + Capacity), returning reciprocal pairs as `cliques()` requires. + +Comparators are Splink/DuckDB SQL: Jaro-Winkler thresholds for names, exact +match for categoricals, a min/max ratio ladder for capacity, and a haversine +distance ladder for position. A 0.5 prior reproduces DUKE's belief update. Only +the *relative* Bayes factor between levels affects ranking, so each comparison's +absolute scale is irrelevant and the operating point is set by one calibrated +`match_probability` threshold per mode (`LINKAGE_THRESHOLD`, `DEDUP_THRESHOLD`). + +## Ground truth + +The production `powerplants.csv` records which source IDs ended in the same +final cluster. From its `projectID` column we recover **567 GEO↔GPD pairs** as a +cross-source linkage ground truth and the intra-GPD pairs for dedup. The harness +is `analysis/benchmark_splink.py`; the DUKE baseline must be run on `master` +(its binaries are removed on this branch). + +## Results — record linkage (GEO × GPD, 12 countries, vs ground truth) + +| backend | precision | recall | F1 | time | +|---|---|---|---|---| +| DUKE | 0.790 | 0.771 | 0.780 | 6.7 s | +| splink (th 0.70, default) | 0.778 | 0.718 | 0.747 | 2.6 s | +| splink (th 0.50) | 0.700 | **0.802** | 0.748 | 2.6 s | +| splink (th 0.90) | **0.806** | 0.681 | 0.738 | 2.7 s | + +Splink lands within ~3 F1-points of DUKE while running **~2.6× faster**, and the +threshold is a single dial trading precision for recall: at `th=0.50` its recall +(0.80) **exceeds** DUKE's (0.77). The matcher's raw recall before the 1:1 +`best_matches` reduction is **0.82**, so the headroom is in candidate ranking, +not in the comparators. As in #301, the absolute F1 (~0.4–0.8 depending on +scope) is a property of the *evaluation* — raw pairwise GEO↔GPD against a +ground truth produced by the full multi-source pipeline. + +## Results — deduplication (GPD, 8 countries) + +| backend | merges | principled | time | +|---|---|---|---| +| DUKE | 454 | — | 21.7 s | +| splink (th 0.95, default) | 459 | 86 % | 3.2 s | +| splink (th 0.99) | 190 | 100 % | 2.9 s | + +A scoreboard against the intra-GPD clusters in `powerplants.csv` is **circular** — +production used DUKE for dedup, so the clusters are DUKE's own output and its +recall is 1.0 by construction. The fair test is what each backend merges: at its +default threshold Splink makes a near-identical number of merges (459 vs DUKE's +454) and **86 %** of the merged pairs share a fueltype *and* lie within 5 km +*and* have a high name similarity — i.e. they are principled. Many are +multi-unit plants (*Altahullion* / *Altahullion Extension*, *Abbots Ripton* / +*Abbots Ripton Solar Farm Ext*), exactly what `aggregate_units` exists to +collapse. Splink ran **~7× faster** than DUKE. + +## Outcome + +DUKE — its Java dependency, bundled `.jar` binaries, and XML configs — was +removed. `linkage` (Splink + DuckDB) is the sole matching engine, selected +automatically with no configuration. `splink` is a hard dependency; Java is no +longer required. + +> Note: a fixed-parameter model (m/u set from the DUKE weights) was chosen over +> Splink's usual unsupervised EM training because the pipeline calls `match` +> once per country, and per-call EM on tiny slices is both slow and unstable. +> Fixing the parameters keeps the engine deterministic and fast while still +> using Splink's vectorised DuckDB Fellegi-Sunter scoring. diff --git a/docs/api-core.md b/docs/api-core.md index 8f13f3f7..8b2dcc19 100644 --- a/docs/api-core.md +++ b/docs/api-core.md @@ -8,6 +8,6 @@ SPDX-License-Identifier: MIT ::: powerplantmatching.core -::: powerplantmatching.duke +::: powerplantmatching.linkage ::: powerplantmatching.accessor diff --git a/docs/license.md b/docs/license.md index d88059b1..1b4784d8 100644 --- a/docs/license.md +++ b/docs/license.md @@ -8,7 +8,6 @@ SPDX-License-Identifier: MIT `powerplantmatching` is released as free software under the [MIT](LICENSES/MIT.txt) license. The default output data [powerplants.csv](powerplants.csv) generated by the package is released under [CC BY 4.0](LICENSES/CC-BY-4.0.txt). -Parts of the repository may be licensed under different licenses, especially dependent package binaries for `duke` being licensed under [Apache 2.0 license](https://github.com/PyPSA/powerplantmatching/tree/master/LICENSES/Apache-2.0.txt). This repository uses the [REUSE](https://reuse.software/) conventions to indicate the licenses that apply to individual files and parts of the repository. For details on the licenses that apply, see the the header information of the respective files and [REUSE.toml](REUSE.toml) for details. diff --git a/docs/release-notes.md b/docs/release-notes.md index 8a405895..b0518281 100644 --- a/docs/release-notes.md +++ b/docs/release-notes.md @@ -8,6 +8,7 @@ SPDX-License-Identifier: MIT ## Upcoming Version +* Replaced the Java-based DUKE matching engine with a Splink (Fellegi-Sunter) record-linkage backend (`powerplantmatching.linkage`, built on `splink` + DuckDB). Matching no longer requires a Java installation or the bundled DUKE binaries. The `parallel_duke_processes` config key is renamed to `parallel_processes`. * OSM dataset upgraded from a Europe-only snapshot (`osm_europe.csv`) to a global snapshot (`osm_global.csv.gz` taken from [`osm-powerplants`](https://github.com/open-energy-transition/osm-powerplants). ## [v0.8.1](https:://github.com/PyPSA/powerplantmatching/releases/tag/v0.8.1) (11th February 2026) diff --git a/powerplantmatching/accessor.py b/powerplantmatching/accessor.py index 05d50145..d5b9a78a 100644 --- a/powerplantmatching/accessor.py +++ b/powerplantmatching/accessor.py @@ -88,12 +88,12 @@ def set_name(self, name): def get_name(self): return self._obj.columns.name - def match_with(self, df, labels=None, config=None, reduced=True, **dukeargs): + def match_with(self, df, labels=None, config=None, reduced=True, **kwargs): from .matching import combine_multiple_datasets, reduce_matched_dataframe from .utils import to_list_if_other dfs = [self._obj] + to_list_if_other(df) - res = combine_multiple_datasets(dfs, labels, config=config, **dukeargs) + res = combine_multiple_datasets(dfs, labels, config=config, **kwargs) if reduced: return res.pipe(reduce_matched_dataframe, config=config) return res diff --git a/powerplantmatching/cleaning.py b/powerplantmatching/cleaning.py index c4df7961..d1fc8a22 100644 --- a/powerplantmatching/cleaning.py +++ b/powerplantmatching/cleaning.py @@ -15,7 +15,7 @@ from deprecation import deprecated from .core import get_config, get_obj_if_Acc -from .duke import duke +from .linkage import match from .utils import get_name, set_column_name logger = logging.getLogger(__name__) @@ -504,11 +504,11 @@ def aggregate_units( country_query = "Country == @c" query = " and ".join(filter(None, [agg_query, block_query, country_query])) duplicates = pd.concat( - [duke(df.query(query), threads=threads) for c in countries] + [match(df.query(query), threads=threads) for c in countries] ) else: query = " and ".join(filter(None, [agg_query, block_query])) - duplicates = duke(df.query(query) if query else df, threads=threads) + duplicates = match(df.query(query) if query else df, threads=threads) df = cliques(df, duplicates) df = df.groupby("grouped").agg(props_for_groups) diff --git a/powerplantmatching/collection.py b/powerplantmatching/collection.py index 8169affa..fef1a167 100644 --- a/powerplantmatching/collection.py +++ b/powerplantmatching/collection.py @@ -31,7 +31,7 @@ def collect( update=False, reduced=True, config=None, - **dukeargs, + **kwargs, ): """ Return the collection for a given list of datasets in matched or @@ -47,7 +47,7 @@ def collect( Switch as to return the reduced (True) or matched (False) dataset. config : dict Configuration file of powerplantmatching - **dukeargs : keyword-args for duke + **kwargs : keyword-args for the matcher """ from . import data @@ -88,7 +88,7 @@ def df_by_name(name): if update: dfs = parmap(df_by_name, datasets) - matched = combine_multiple_datasets(dfs, datasets, config=config, **dukeargs) + matched = combine_multiple_datasets(dfs, datasets, config=config, **kwargs) ( matched.assign(projectID=lambda df: df.projectID.astype(str)).to_csv( outfn_matched, index_label="id" diff --git a/powerplantmatching/duke.py b/powerplantmatching/duke.py deleted file mode 100644 index 4c6fff7a..00000000 --- a/powerplantmatching/duke.py +++ /dev/null @@ -1,160 +0,0 @@ -# SPDX-FileCopyrightText: Contributors to powerplantmatching -# -# SPDX-License-Identifier: MIT - -import logging -import os -import shutil -import subprocess as sub -import tempfile - -import numpy as np -import pandas as pd - -from .core import _package_data - -logger = logging.getLogger(__name__) - - -def add_geoposition_for_duke(df): - """ - Returns the same pandas.Dataframe with an additional column "Geoposition" - which concats the latitude and longitude of the powerplant in a string - - """ - if not df.loc[:, ["lat", "lon"]].isnull().all().all(): - return df.assign( - Geoposition=df[["lat", "lon"]] - .astype(str) - .apply(lambda s: ",".join(s), axis=1) - .replace("nan,nan", np.nan) - ) - else: - return df.assign(Geoposition=np.nan) - - -def duke( - datasets, - labels=["one", "two"], - singlematch=False, - showmatches=False, - keepfiles=False, - showoutput=False, - threads=1, -): - """ - Run duke in different modes (Deduplication or Record Linkage Mode) to - either locate duplicates in one database or find the similar entries in two - different datasets. In RecordLinkagesMode (match two databases) please - set singlematch=True and use best_matches() afterwards - - Parameters - ---------- - - datasets : pd.DataFrame or [pd.DataFrame] - A single dataframe is run in deduplication mode, while multiple ones - are linked - labels : [str], default ['one', 'two'] - Labels for the linked dataframe - singlematch: boolean, default False - Only in Record Linkage Mode. Only report the best match for each entry - of the first named dataset. This does not guarantee a unique match in - the second named dataset. - keepfiles : boolean, default False - If true, do not delete temporary files - """ - - try: - sub.run(["java", "-version"], check=True, capture_output=True) - except sub.CalledProcessError: - err = "Java is not installed or not in the system's PATH. Please install Java and ensure it is in your system's PATH, then try again." - logger.error(err) - raise FileNotFoundError(err) - - dedup = isinstance(datasets, pd.DataFrame) - if dedup: - # Deduplication mode - duke_config = "Deleteduplicates.xml" - datasets = [datasets] - else: - duke_config = "Comparison.xml" - - duke_bin_dir = _package_data("duke_binaries") - - os.environ["CLASSPATH"] = os.pathsep.join( - [os.path.join(duke_bin_dir, r) for r in os.listdir(duke_bin_dir)] - ) - tmpdir = tempfile.mkdtemp() - - try: - shutil.copyfile( - os.path.join(_package_data(duke_config)), os.path.join(tmpdir, "config.xml") - ) - - logger.debug("Comparing files: %s", ", ".join(labels)) - - for n, df in enumerate(datasets): - df = add_geoposition_for_duke(df) - # due to index unity (see https://github.com/larsga/Duke/issues/236) - if n == 1: - shift_by = datasets[0].index.max() + 1 - df.index += shift_by - df.to_csv(os.path.join(tmpdir, f"file{n + 1}.csv"), index_label="id") - if n == 1: - df.index -= shift_by - - args = [ - "java", - "-Dfile.encoding=UTF-8", - "no.priv.garshol.duke.Duke", - "--linkfile=linkfile.txt", - f"--threads={threads}", - ] - if singlematch: - args.append("--singlematch") - if showmatches: - args.append("--showmatches") - stdout = sub.PIPE - else: - stdout = None - args.append("config.xml") - - run = sub.Popen( - args, - stderr=sub.PIPE, - cwd=tmpdir, - stdout=stdout, - universal_newlines=True, - ) - _, stderr = run.communicate() - - if showmatches: - print(_) - - logger.debug(f"Stderr: {stderr}") - if any(word in stderr.lower() for word in ["error", "fehler"]): - raise RuntimeError(f"duke failed: {stderr}") - - if dedup: - return pd.read_csv( - os.path.join(tmpdir, "linkfile.txt"), - encoding="utf-8", - usecols=[1, 2], - names=labels, - ) - else: - res = pd.read_csv( - os.path.join(tmpdir, "linkfile.txt"), - usecols=[1, 2, 3], - names=labels + ["scores"], - ) - res.iloc[:, 1] -= shift_by - res["scores"] = res.scores.astype(float) - return res - - finally: - if keepfiles: - logger.debug(f"Files of the duke run are kept in {tmpdir}") - else: - shutil.rmtree(tmpdir) - logger.debug(f"Files of the duke run have been deleted in {tmpdir}") diff --git a/powerplantmatching/linkage.py b/powerplantmatching/linkage.py new file mode 100644 index 00000000..898caeb5 --- /dev/null +++ b/powerplantmatching/linkage.py @@ -0,0 +1,218 @@ +# SPDX-FileCopyrightText: Contributors to powerplantmatching +# +# SPDX-License-Identifier: MIT + +""" +Splink-backed record-linkage and deduplication engine. + +``match`` takes a list of two frames for record linkage or a single frame for +deduplication and returns the matched index pairs. It runs a fixed-parameter +`Splink `_ (Fellegi-Sunter) +model on DuckDB: the per-field ``m``/``u`` probabilities are set directly from +the former DUKE field weights (``Comparison.xml`` / ``Deleteduplicates.xml``), +so the model needs no per-call EM training and is fully deterministic. + +Comparators: Jaro-Winkler thresholds for names, exact match for categorical +fields, a min/max ratio for capacity and a haversine distance ladder for +position. A 0.5 prior reproduces DUKE's belief update; only the relative +Bayes factor between levels matters for ranking, so the absolute match +probability is calibrated through ``LINKAGE_THRESHOLD`` / ``DEDUP_THRESHOLD``. +""" + +import logging +from dataclasses import dataclass + +import numpy as np +import pandas as pd +from splink import DuckDBAPI, Linker + +logger = logging.getLogger(__name__) +logging.getLogger("splink").setLevel(logging.ERROR) + +_HAVERSINE_KM = ( + '2*6371*asin(sqrt(pow(sin(radians(("lat_r"-"lat_l"))/2),2)' + '+cos(radians("lat_l"))*cos(radians("lat_r"))' + '*pow(sin(radians(("lon_r"-"lon_l"))/2),2)))' +) + + +@dataclass +class FieldSpec: + """A field comparison expressed as similarity levels with DUKE-derived probabilities. + + ``levels`` maps each SQL agreement condition to the DUKE agreement + probability ``p`` for records that reach it (high for strong agreement, low + for the catch-all ``ELSE`` level). + """ + + output: str + null_sql: str + levels: list # [(sql_condition, p)], last entry is the ELSE level + + def comparison(self) -> dict: + sqls = [sql for sql, _ in self.levels] + ps = np.array([p for _, p in self.levels], dtype=float) + m = ps / ps.sum() + u = (1 - ps) / (1 - ps).sum() + graded = [ + { + "sql_condition": sql, + "label_for_charts": sql, + "m_probability": float(mi), + "u_probability": float(ui), + } + for sql, mi, ui in zip(sqls, m, u) + ] + null_level = { + "sql_condition": self.null_sql, + "label_for_charts": "null", + "is_null_level": True, + } + return { + "output_column_name": self.output, + "comparison_levels": [null_level, *graded], + } + + +def _name(low: float, high: float) -> FieldSpec: + return FieldSpec( + "Name", + '"Name_l" IS NULL OR "Name_r" IS NULL', + [ + ('jaro_winkler_similarity("Name_l","Name_r") >= 0.95', high), + ( + 'jaro_winkler_similarity("Name_l","Name_r") >= 0.85', + low + 0.7 * (high - low), + ), + ( + 'jaro_winkler_similarity("Name_l","Name_r") >= 0.7', + low + 0.35 * (high - low), + ), + ("ELSE", low), + ], + ) + + +def _exact(col: str, low: float, high: float) -> FieldSpec: + return FieldSpec( + col, + f'"{col}_l" IS NULL OR "{col}_r" IS NULL', + [(f'"{col}_l" = "{col}_r"', high), ("ELSE", low)], + ) + + +def _capacity(low: float, high: float) -> FieldSpec: + ratio = ( + 'least("Capacity_l","Capacity_r")/nullif(greatest("Capacity_l","Capacity_r"),0)' + ) + return FieldSpec( + "Capacity", + '"Capacity_l" IS NULL OR "Capacity_r" IS NULL', + [ + (f"{ratio} >= 0.95", high), + (f"{ratio} >= 0.8", low + 0.5 * (high - low)), + ("ELSE", low), + ], + ) + + +def _geo(low: float, high: float) -> FieldSpec: + return FieldSpec( + "geo", + '"lat_l" IS NULL OR "lat_r" IS NULL OR "lon_l" IS NULL OR "lon_r" IS NULL', + [ + (f"{_HAVERSINE_KM} <= 1", high), + (f"{_HAVERSINE_KM} <= 5", low + 0.7 * (high - low)), + ("ELSE", low), + ], + ) + + +LINKAGE_FIELDS = [ + _name(0.09, 0.99), + _exact("Fueltype", 0.09, 0.7), + _exact("Country", 0.01, 0.53), + _capacity(0.3, 0.75), + _geo(0.1, 0.8), +] +LINKAGE_THRESHOLD = 0.7 + +DEDUP_FIELDS = [ + _name(0.09, 0.99), + _exact("Fueltype", 0.05, 0.65), + _exact("Technology", 0.25, 0.51), + _exact("Country", 0.05, 0.51), + _capacity(0.49, 0.51), + _geo(0.05, 0.75), +] +DEDUP_THRESHOLD = 0.95 + + +def _settings(link_type: str, fields) -> dict: + return { + "link_type": link_type, + "probability_two_random_records_match": 0.5, + "unique_id_column_name": "id", + "blocking_rules_to_generate_predictions": [{"blocking_rule": "1=1"}], + "comparisons": [f.comparison() for f in fields], + "retain_matching_columns": False, + "retain_intermediate_calculation_columns": False, + } + + +def _predict(tables, link_type, fields, threshold) -> pd.DataFrame: + aliases = ["df_left", "df_right"] if isinstance(tables, list) else None + linker = Linker( + tables, _settings(link_type, fields), DuckDBAPI(), input_table_aliases=aliases + ) + return linker.inference.predict( + threshold_match_probability=threshold + ).as_pandas_dataframe() + + +def _deduplicate(df: pd.DataFrame, labels, threshold: float) -> pd.DataFrame: + if len(df) < 2: + return pd.DataFrame(columns=labels) + table = df.assign(id=df.index) + pred = _predict(table, "dedupe_only", DEDUP_FIELDS, threshold) + if pred.empty: + return pd.DataFrame(columns=labels) + a, b = pred["id_l"].to_numpy(), pred["id_r"].to_numpy() + return pd.DataFrame( + {labels[0]: np.concatenate([a, b]), labels[1]: np.concatenate([b, a])} + ) + + +def match(datasets, labels=["one", "two"], singlematch=False, threshold=None, **_): + """ + Link two datasets or deduplicate one. + + A single frame is deduplicated (returns reciprocal index pairs as + ``cliques()`` requires); a list of two frames is linked. In record linkage + pass ``singlematch=True`` and reduce with ``best_matches()`` afterwards. + """ + if isinstance(datasets, pd.DataFrame): + th = DEDUP_THRESHOLD if threshold is None else threshold + return _deduplicate(datasets, labels, th) + + left, right = datasets + if left.empty or right.empty: + return pd.DataFrame(columns=[*labels, "scores"]) + + th = LINKAGE_THRESHOLD if threshold is None else threshold + tables = [left.assign(id=left.index), right.assign(id=right.index)] + pred = _predict(tables, "link_only", LINKAGE_FIELDS, th) + if pred.empty: + return pd.DataFrame(columns=[*labels, "scores"]) + + from_left = pred["source_dataset_l"].to_numpy() == "df_left" + res = pd.DataFrame( + { + labels[0]: np.where(from_left, pred["id_l"], pred["id_r"]), + labels[1]: np.where(from_left, pred["id_r"], pred["id_l"]), + "scores": pred["match_probability"].to_numpy(), + } + ) + if singlematch and not res.empty: + res = res.loc[res.groupby(labels[0])["scores"].idxmax()].reset_index(drop=True) + return res.reset_index(drop=True) diff --git a/powerplantmatching/matching.py b/powerplantmatching/matching.py index 5e9b5690..a4957fe7 100644 --- a/powerplantmatching/matching.py +++ b/powerplantmatching/matching.py @@ -14,7 +14,7 @@ from .cleaning import clean_technology from .core import get_config, get_obj_if_Acc -from .duke import duke +from .linkage import match from .utils import get_name, parmap, read_csv_if_string logger = logging.getLogger(__name__) @@ -22,13 +22,13 @@ def best_matches(links): """ - Subsequent to duke() with singlematch=True. Returns reduced list of + Subsequent to match() with singlematch=True. Returns reduced list of matches on the base of the highest score for each duplicated entry. Parameters ---------- links : pd.DataFrame - Links as returned by duke + Links as returned by match """ labels = links.columns.difference({"scores"}) if links.empty: @@ -39,9 +39,9 @@ def best_matches(links): ) -def compare_two_datasets(dfs, labels, country_wise=True, config=None, **dukeargs): +def compare_two_datasets(dfs, labels, country_wise=True, config=None, **kwargs): """ - Duke-based horizontal match of two databases. Returns the matched + Horizontal match of two databases. Returns the matched dataframe including only the matched entries in a multi-indexed pandas.Dataframe. Compares all properties of the given columns ['Name','Fueltype', 'Technology', 'Country', @@ -49,9 +49,7 @@ def compare_two_datasets(dfs, labels, country_wise=True, config=None, **dukeargs powerplant in different two datasets. The match is in one-to-one mode, that is every entry of the initial databases has maximally one link in order to obtain unique entries in the resulting - dataframe. Attention: When aborting this command, the duke - process will still continue in the background, wait until the - process is finished before restarting. + dataframe. Parameters ---------- @@ -66,24 +64,24 @@ def compare_two_datasets(dfs, labels, country_wise=True, config=None, **dukeargs config = get_config() deprecated_args = {"use_saved_matches", "use_saved_aggregation"} - used_deprecated_args = deprecated_args.intersection(dukeargs) + used_deprecated_args = deprecated_args.intersection(kwargs) if used_deprecated_args: for arg in used_deprecated_args: - dukeargs.pop(arg) + kwargs.pop(arg) msg = "The following arguments were deprecated and are being ignored: " logger.warning(msg + f"{used_deprecated_args}") dfs = list(map(read_csv_if_string, dfs)) - if "singlematch" not in dukeargs: - dukeargs["singlematch"] = True + if "singlematch" not in kwargs: + kwargs["singlematch"] = True def country_link(dfs, country): # country_selector for both dataframes sel_country_b = [df["Country"] == country for df in dfs] # only append if country appears in both dataframse if all(sel.any() for sel in sel_country_b): - return duke( - [df[sel] for df, sel in zip(dfs, sel_country_b)], labels, **dukeargs + return match( + [df[sel] for df, sel in zip(dfs, sel_country_b)], labels, **kwargs ) else: return pd.DataFrame(columns=[*labels, "scores"]) @@ -97,7 +95,7 @@ def country_link(dfs, country): else: links = pd.DataFrame(columns=[*labels, "scores"]) else: - links = duke(dfs, labels=labels, **dukeargs) + links = match(dfs, labels=labels, **kwargs) if links.empty: matches = pd.DataFrame(columns=labels) @@ -118,8 +116,7 @@ def cross_matches(sets_of_pairs, labels=None): ---------- sets_of_pairs : list list of pd.Dataframe's containing only the matches (without - scores), obtained from the linkfile (duke() and - best_matches()) + scores), obtained from match() and best_matches() labels : list of strings list of names of the databases, used for specifying the order of the output @@ -167,10 +164,10 @@ def cross_matches(sets_of_pairs, labels=None): def link_multiple_datasets( - datasets, labels, use_saved_matches=False, config=None, **dukeargs + datasets, labels, use_saved_matches=False, config=None, **kwargs ): """ - Duke-based horizontal match of multiple databases. Returns the + Horizontal match of multiple databases. Returns the matching indices of the datasets. Compares all properties of the given columns ['Name','Fueltype', 'Technology', 'Country', 'Capacity','lat', 'lon'] in order to determine the same @@ -197,7 +194,7 @@ def link_multiple_datasets( def comp_dfs(dfs_lbs): logger.info("Comparing data sources `{}` and `{}`".format(*dfs_lbs[2:])) - return compare_two_datasets(dfs_lbs[:2], dfs_lbs[2:], config=config, **dukeargs) + return compare_two_datasets(dfs_lbs[:2], dfs_lbs[2:], config=config, **kwargs) mapargs = [[dfs[c], dfs[d], labels[c], labels[d]] for c, d in combs] all_matches = parmap(comp_dfs, mapargs) @@ -205,9 +202,9 @@ def comp_dfs(dfs_lbs): return cross_matches(all_matches, labels=labels) -def combine_multiple_datasets(datasets, labels=None, config=None, **dukeargs): +def combine_multiple_datasets(datasets, labels=None, config=None, **kwargs): """ - Duke-based horizontal match of multiple databases. Returns the + Horizontal match of multiple databases. Returns the matched dataframe including only the matched entries in a multi-indexed pandas.Dataframe. Compares all properties of the given columns ['Name','Fueltype', 'Technology', 'Country', @@ -252,7 +249,7 @@ def combined_dataframe(cross_matches, datasets, config): .reset_index(drop=True) ) - crossmatches = link_multiple_datasets(datasets, labels, config=config, **dukeargs) + crossmatches = link_multiple_datasets(datasets, labels, config=config, **kwargs) return combined_dataframe(crossmatches, datasets, config).reindex( columns=config["target_columns"], level=0 ) diff --git a/powerplantmatching/package_data/Comparison.xml b/powerplantmatching/package_data/Comparison.xml deleted file mode 100644 index cfa11fbe..00000000 --- a/powerplantmatching/package_data/Comparison.xml +++ /dev/null @@ -1,100 +0,0 @@ - - - - - - - - - - - 0.965 - - - ID - - - NAME - no.priv.garshol.duke.comparators.JaroWinklerTokenized - 0.09 - 0.99 - - - FUELTYPE - no.priv.garshol.duke.comparators.QGramComparator - 0.09 - 0.7 - - - COUNTRY - no.priv.garshol.duke.comparators.QGramComparator - 0.0 - 0.53 - - - CAPACITY - no.priv.garshol.duke.comparators.NumericComparator - 0.3 - 0.75 - - - GEOPOSITION - GeoComparator - 0.1 - 0.8 - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - diff --git a/powerplantmatching/package_data/Deleteduplicates.xml b/powerplantmatching/package_data/Deleteduplicates.xml deleted file mode 100644 index b7490fa8..00000000 --- a/powerplantmatching/package_data/Deleteduplicates.xml +++ /dev/null @@ -1,89 +0,0 @@ - - - - - - - - - - - 0.96 - - - ID - - - NAME - no.priv.garshol.duke.comparators.JaroWinklerTokenized - 0.09 - 0.99 - - - FUELTYPE - no.priv.garshol.duke.comparators.QGramComparator - 0.05 - 0.65 - - - TECHNOLOGY - no.priv.garshol.duke.comparators.QGramComparator - 0.25 - 0.51 - - - COUNTRY - no.priv.garshol.duke.comparators.QGramComparator - 0.05 - 0.51 - - - CAPACITY - no.priv.garshol.duke.comparators.NumericComparator - 0.49 - 0.51 - - - GEOPOSITION - GeoComparator - 0.05 - 0.75 - - - - - - - - - - - - - - - - - - - - - - diff --git a/powerplantmatching/package_data/config.yaml b/powerplantmatching/package_data/config.yaml index aced4f37..8afe6c72 100644 --- a/powerplantmatching/package_data/config.yaml +++ b/powerplantmatching/package_data/config.yaml @@ -53,7 +53,7 @@ fully_included_sources: aggregate_only_matching_sources: - MASTR # the matching process of very small units is not efficient -parallel_duke_processes: false +parallel_processes: false threads_extend_by_non_matched: 16 matched_data_url: https://raw.githubusercontent.com/PyPSA/powerplantmatching/{tag}/powerplants.csv diff --git a/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT-tests.jar b/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT-tests.jar deleted file mode 100644 index d39b8e78..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT-tests.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT.jar deleted file mode 100644 index b2ce083e..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-core-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-es-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-es-1.3-SNAPSHOT.jar deleted file mode 100644 index bb237b8d..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-es-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-json-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-json-1.3-SNAPSHOT.jar deleted file mode 100644 index dbd239ad..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-json-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-lucene-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-lucene-1.3-SNAPSHOT.jar deleted file mode 100644 index 3bbb42d9..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-lucene-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-mapdb-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-mapdb-1.3-SNAPSHOT.jar deleted file mode 100644 index ebd8104a..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-mapdb-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-mongodb-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-mongodb-1.3-SNAPSHOT.jar deleted file mode 100644 index 9838fb62..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-mongodb-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/package_data/duke_binaries/duke-server-1.3-SNAPSHOT.jar b/powerplantmatching/package_data/duke_binaries/duke-server-1.3-SNAPSHOT.jar deleted file mode 100644 index b294adee..00000000 Binary files a/powerplantmatching/package_data/duke_binaries/duke-server-1.3-SNAPSHOT.jar and /dev/null differ diff --git a/powerplantmatching/utils.py b/powerplantmatching/utils.py index 4b9396bc..8840c5a6 100644 --- a/powerplantmatching/utils.py +++ b/powerplantmatching/utils.py @@ -341,7 +341,7 @@ def parmap(f, arg_list, config=None, threads=None): """ Parallel mapping function. Use this function to parallelly map function f onto arguments in arg_list. The maximum number of parallel threads is - taken from config.yaml:parallel_duke_processes. + taken from config.yaml:parallel_processes. Parameters --------- @@ -359,7 +359,7 @@ def parmap(f, arg_list, config=None, threads=None): config = get_config() if threads is None: - threads = config["parallel_duke_processes"] + threads = config["parallel_processes"] if isinstance(threads, bool): threads = config.get("process_limit", 1) diff --git a/pyproject.toml b/pyproject.toml index f8626379..719f7687 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -49,6 +49,7 @@ dependencies = [ "deprecation", "tqdm", "openpyxl", + "splink>=4", ] [project.urls] diff --git a/test/test_duke.py b/test/test_duke.py deleted file mode 100644 index d1ea14d5..00000000 --- a/test/test_duke.py +++ /dev/null @@ -1,97 +0,0 @@ -# SPDX-FileCopyrightText: Contributors to powerplantmatching -# -# SPDX-License-Identifier: MIT - -import numpy as np -import pandas as pd - -from powerplantmatching.duke import add_geoposition_for_duke, duke - -TEST_DATA = { - "Name": [ - "Powerplant", - "an hydro powerplant", - " another powerplant with whitespaces", - " Power II coalition", - " Kraftwerk Besonders besonders '2' CHP", - ], - "Fueltype": [ - "", - "Run of-River", - "OCGT", - "Nuclear Power", - "", - ], - "Technology": [ - "Natural Gas", - "Run of-River", - "", - " Nuclear", - "", - ], - "Set": [ - np.nan, - "", - "", - "Powerplant", - "", - ], - "Country": ["", "", "", "", ""], - "Capacity": [120, 40, 80.008, 90.122, 400.2], - "lat": [51.5074, np.nan, 40.7128, 42.8765, 48.2031], - "lon": [-0.1278, np.nan, -74.0060, -60.9807, -1.0221], -} - - -def test_add_geoposition_for_duke(): - df = pd.DataFrame(TEST_DATA) - - expected = pd.DataFrame(TEST_DATA) - expected["Geoposition"] = [ - "51.5074,-0.1278", - np.nan, - "40.7128,-74.006", - "42.8765,-60.9807", - "48.2031,-1.0221", - ] - - output = add_geoposition_for_duke(df) - - pd.testing.assert_frame_equal(output, expected) - - -def test_duke_deduplication_mode(): - df = pd.DataFrame(TEST_DATA) - - expected = pd.DataFrame({"one": [0, 1], "two": [1, 0]}) - - output = duke(df) - - pd.testing.assert_frame_equal(output, expected) - - -def test_duke_record_linkage_mode(): - df1 = pd.DataFrame(TEST_DATA) - - df2 = pd.DataFrame(TEST_DATA) - df2.loc[0, "Name"] = "Plant" - df2.loc[0, "Technology"] = "Gas" - df2.loc[2, "Capacity"] = 44.5 - - expected = pd.DataFrame( - { - "one": {0: 0, 1: 1, 2: 2, 3: 3, 4: 4}, - "two": {0: 1, 1: 1, 2: 2, 3: 3, 4: 4}, - "scores": { - 0: 0.9769736842105264, - 1: 0.9985590778097984, - 2: 0.9992083247849796, - 3: 0.999639379733141, - 4: 0.9991589571068124, - }, - } - ) - - output = duke([df1, df2], labels=["one", "two"], singlematch=False) - - pd.testing.assert_frame_equal(output, expected) diff --git a/test/test_linkage.py b/test/test_linkage.py new file mode 100644 index 00000000..233e35d1 --- /dev/null +++ b/test/test_linkage.py @@ -0,0 +1,79 @@ +# SPDX-FileCopyrightText: Contributors to powerplantmatching +# +# SPDX-License-Identifier: MIT + +import pandas as pd +import pytest + +from powerplantmatching import linkage as lk + +TEST_DATA = { + "Name": ["Powerplant", "Hydro Station", "Gas Turbine A", "Coal Block 1"], + "Fueltype": ["Hard Coal", "Hydro", "Natural Gas", "Hard Coal"], + "Technology": ["Steam Turbine", "Run-Of-River", "CCGT", "Steam Turbine"], + "Country": ["Germany", "Germany", "Germany", "Germany"], + "Capacity": [120.0, 40.0, 80.0, 300.0], + "lat": [51.5, 48.1, 50.9, 52.5], + "lon": [7.0, 11.6, 6.9, 13.4], +} + + +@pytest.fixture +def left(): + return pd.DataFrame(TEST_DATA) + + +@pytest.fixture +def right(left): + df = left.copy() + df.loc[0, "Name"] = "Power Plant" # near-duplicate of "Powerplant" + df.loc[2, "Capacity"] = 82.0 + return df + + +def test_record_linkage_format(left, right): + out = lk.match([left, right], labels=["one", "two"], singlematch=True) + assert list(out.columns) == ["one", "two", "scores"] + assert out["scores"].between(0, 1).all() + assert (out["scores"] >= lk.LINKAGE_THRESHOLD).all() + + +def test_record_linkage_matches_identical_rows(left, right): + out = lk.match([left, right], labels=["one", "two"], singlematch=True) + pairs = set(zip(out["one"], out["two"])) + assert {(i, i) for i in left.index}.issubset(pairs) + + +def test_singlematch_is_unique_per_left(left, right): + out = lk.match([left, right], labels=["one", "two"], singlematch=True) + assert out["one"].is_unique + + +def test_empty_input_returns_empty_links(left): + out = lk.match([left, left.iloc[0:0]], labels=["one", "two"]) + assert out.empty + assert list(out.columns) == ["one", "two", "scores"] + + +def test_dedup_returns_symmetric_pairs(left): + dup = left.copy() + dup.loc[len(dup)] = dup.loc[0] # exact duplicate of row 0 + out = lk.match(dup, labels=["one", "two"]) + assert list(out.columns) == ["one", "two"] + forward = set(zip(out["one"], out["two"])) + assert forward == set( + zip(out["two"], out["one"]) + ) # reciprocal, as cliques() requires + assert (0, len(dup) - 1) in forward + + +def test_geo_proximity_increases_score(left): + one = left.iloc[[0]].reset_index(drop=True) + other = one.copy() + other.loc[0, "Name"] = "Power Station" # only loosely similar, so geo matters + near, far = other.copy(), other.copy() + near.loc[0, ["lat", "lon"]] = left.loc[0, ["lat", "lon"]].to_numpy() + far.loc[0, ["lat", "lon"]] = [40.0, -3.0] # far from row 0 + + score = lambda r: lk.match([one, r], threshold=0.0)["scores"].iloc[0] + assert score(near) > score(far)