Transaction matching across two parties¶
About this example¶
This example demonstrates a two-step pipeline for transaction matching inside a Data Clean Room (DCR). Two parties bring their transaction datasets and identify matching records without exposing raw data to each other.
The challenge at scale is memory: loading billions of rows into pandas exhausts node memory before any matching logic runs. The pipeline addresses this by splitting work across two layers:
- SQL pre-filter (Step 1): A SQL template runs an exact-key join on the Snowflake warehouse. The hash join reduces billions of rows to a manageable candidate set in seconds with no memory constraints.
- Ray task matching (Step 2): An ML Jobs template reads the reduced candidate set and distributes custom matching logic across Ray workers on a compute pool. Because the candidate set fits in memory, the Python script can run any matching algorithm: fuzzy scoring, clustering, entity resolution, or ML model inference.
Other applications of this approach:
- Financial reconciliation: Match payments and settlements across banks or payment processors on shared customer identifiers.
- Fraud detection: Identify overlapping suspicious transactions across institutions without sharing customer records.
- Retail media measurement: Join ad exposure logs with purchase transactions to measure campaign lift.
- Insurance claims matching: Correlate claims data across carriers and providers for duplicate detection.
Roles:
- Owner (collaboration creator and analysis runner): Registers the owner transaction data, stages the matching script, creates the collaboration, provisions the compute pool, and runs both pipeline steps.
- Collaborator: Registers the collaborator transaction data, reviews and joins the collaboration, and links the data offering.
Pipeline:
- SQL pre-filter: Join both transaction tables on exact keys to produce a
cleanroom.match_candidatestable containing only candidate pairs. - Ray matching: Load the candidate set from
cleanroom.match_candidatesand distribute the custom matching function across Ray workers on a multi-node compute pool.
Prerequisites¶
-
Two accounts with the Data Clean Rooms environment installed. For cross-region deployments, enable Cross-Cloud Auto-Fulfillment.
-
Both parties’ transaction tables must use SHA-256 hashed join keys. Raw PII must be hashed before registration.
-
The analysis runner account must have a compute pool available.
CPU_X64_Lis recommended for candidate sets in the tens of millions of rows. -
Generate sample data by running the sample data generator notebook in both accounts. Upload the notebook to Snowsight (Notebooks » Import .ipynb file), set the
DATABASE_NAMEandSCHEMA_NAMEvariables in the first code cell, then run all cells. The notebook creates:PROVIDER_TRANSACTIONSin the owner accountPARTNER_TRANSACTIONSin the collaborator account
-
Upload the matching script (txn_match.py) to a stage in the owner account:
Run the example¶
Download and run the following SQL worksheets in the owner and collaborator accounts. The worksheets cover data offering registration, template and code spec registration, collaboration creation, compute pool setup, and execution of both pipeline steps.
- Owner worksheet: Run this in the owner account.
- Collaborator worksheet: Run this in the collaborator account.
Step 1 runs on the warehouse and completes in seconds:
Step 2 runs on the compute pool and distributes matching across Ray workers:
Customize the matching logic¶
The matching script (txn_match.py) reads from cleanroom.match_candidates,
distributes one Ray task per segment, and writes results to
cleanroom.match_results. The downloadable script includes date-proximity
matching. The simplified version below shows the structure of the function you
replace:
To use a library not already in the container runtime (for example, scikit-learn
or rapidfuzz), add it to pip_requirements in the code spec:
Important
get_active_session() is only available in the main process of the ML Jobs
script. Don’t call it inside @ray.remote worker functions. All data reads and
writes must happen in the main process. Workers receive pandas DataFrames as
arguments and return plain Python dicts.
DCR views rename columns based on schema_and_template_policies. The
join_standard / hashed_email_sha256 policy renames HASHED_EMAIL to
HASHED_EMAIL_SHA256. The timestamp category renames date columns to
TIMESTAMP. Always detect columns dynamically or use the renamed names.