{ "cells": [ { "cell_type": "markdown", "metadata": {}, "source": [ "# Advanced Deduplication Strategies in Polars\n", "\n", "This notebook accompanies the blog post on **Advanced Deduplication Strategies in Polars: Beyond Sort and Drop**.\n", "It provides runnable code examples for basic unique extraction, expression masking, latest-record patterns, conflict resolution, near-duplicate debouncing, and streaming out-of-core deduplication." ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 1. Imports and Dataset Setup" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "import polars as pl\n", "import pandas as pd\n", "\n", "# Sample event dataset with duplicate users and timestamps\n", "data = {\n", " \"user_id\": [101, 102, 101, 103, 102, 101, 104],\n", " \"timestamp\": [100, 105, 102, 101, 108, 99, 103],\n", " \"status\": [\"click\", \"view\", \"purchase\", \"view\", \"click\", \"view\", \"purchase\"],\n", " \"amount\": [10.0, None, 50.0, 5.0, 15.0, 10.0, 100.0]\n", "}\n", "\n", "df = pl.DataFrame(data)\n", "print(df)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 2. Core Deduplication API: `unique()`\n", "Compare preserving physical order (`maintain_order=True`) vs maximum parallel throughput (`keep='any'`, `maintain_order=False`)." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Strategy A: Keep first occurrence preserving physical order\n", "df_first = df.unique(subset=[\"user_id\"], keep=\"first\", maintain_order=True)\n", "print(\"Keep First (maintain_order=True):\")\n", "print(df_first)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Strategy B: Maximum parallel performance across CPU cores\n", "df_fast = df.unique(subset=[\"user_id\"], keep=\"any\", maintain_order=False)\n", "print(\"Keep Any (maintain_order=False):\")\n", "print(df_fast)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 3. Duplicate Inspection & Masking Expressions\n", "Use expression contexts (`.is_duplicated()`, `.is_unique()`, `.is_first_distinct()`, `.is_last_distinct()`) without modifying the DataFrame." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "duplicate_analysis = df.select([\n", " pl.col(\"user_id\"),\n", " pl.col(\"user_id\").is_duplicated().alias(\"is_dup\"),\n", " pl.col(\"user_id\").is_unique().alias(\"is_uniq\"),\n", " pl.col(\"user_id\").is_first_distinct().alias(\"is_first\"),\n", " pl.col(\"user_id\").is_last_distinct().alias(\"is_last\")\n", "])\n", "print(duplicate_analysis)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Filter out all duplicate entries (keep pristine unique records only)\n", "df_pristine = df.filter(pl.col(\"user_id\").is_unique())\n", "print(df_pristine)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 4. \"Latest Record\" Pattern: 4 Polars Alternatives\n", "Extract the record with the latest timestamp per entity." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Method 1: Sort by timestamp descending + unique()\n", "res1 = df.sort(\"timestamp\", descending=True).unique(subset=[\"user_id\"], keep=\"first\", maintain_order=True)\n", "print(\"Method 1 (sort + unique):\")\n", "print(res1)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Method 2: Window expression max().over()\n", "res2 = df.filter(pl.col(\"timestamp\") == pl.col(\"timestamp\").max().over(\"user_id\"))\n", "print(\"Method 2 (max().over()):\")\n", "print(res2)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Method 3: Deterministic tie-breaking with sort_by().first().over()\n", "res3 = df.filter(pl.col(\"timestamp\") == pl.col(\"timestamp\").sort_by(\"timestamp\", descending=True).first().over(\"user_id\"))\n", "print(\"Method 3 (sort_by().first().over()):\")\n", "print(res3)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Method 4: GroupBy aggregation with in-group expression sorting\n", "res4 = df.group_by(\"user_id\").agg([\n", " pl.col(\"timestamp\").sort_by(\"timestamp\").last().alias(\"latest_timestamp\"),\n", " pl.col(\"status\").sort_by(\"timestamp\").last().alias(\"latest_status\"),\n", " pl.col(\"amount\").sort_by(\"timestamp\").last().alias(\"latest_amount\")\n", "])\n", "print(\"Method 4 (group_by + sort_by):\")\n", "print(res4)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 5. Conflict Resolution & Custom Aggregations\n", "Resolve non-key column contradictions without dropping non-null data." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "user_updates = pl.DataFrame({\n", " \"user_id\": [1, 1, 1, 2, 2],\n", " \"email\": [\"alice@old.com\", None, \"alice@new.com\", \"bob@work.com\", None],\n", " \"phone\": [None, \"+1-555-0199\", None, None, \"+1-555-0200\"],\n", " \"tags\": [[\"vip\"], [\"buyer\"], [\"vip\", \"active\"], [\"lead\"], [\"buyer\"]]\n", "})\n", "\n", "resolved_users = user_updates.group_by(\"user_id\").agg([\n", " pl.col(\"email\").drop_nulls().last().alias(\"email\"),\n", " pl.col(\"phone\").drop_nulls().first().alias(\"phone\"),\n", " pl.col(\"tags\").explode().unique().alias(\"all_tags\")\n", "])\n", "print(resolved_users)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 6. Near-Duplicate & Time-Window Event Debouncing\n", "Debounce clickstream events occurring within 5 seconds of the prior click per user." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "click_stream = pl.DataFrame({\n", " \"user_id\": [101, 101, 101, 102, 102],\n", " \"timestamp\": [100, 102, 120, 200, 204],\n", " \"action\": [\"button_click\", \"button_click\", \"button_click\", \"checkout\", \"checkout\"]\n", "})\n", "\n", "debounced = click_stream.with_columns(\n", " time_delta = pl.col(\"timestamp\") - pl.col(\"timestamp\").shift(1).over(\"user_id\")\n", ").filter(\n", " pl.col(\"time_delta\").is_null() | (pl.col(\"time_delta\") > 5)\n", ").drop(\"time_delta\")\n", "\n", "print(debounced)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 7. Lazy Evaluation & Out-of-Core Deduplication\n", "Build a streaming pipeline using LazyFrame.unique()." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "lazy_plan = df.lazy().unique(subset=[\"user_id\"], keep=\"any\", maintain_order=False)\n", "print(\"Execution Plan:\")\n", "print(lazy_plan.explain())\n", "\n", "print(\"Collected Result:\")\n", "print(lazy_plan.collect())" ] } ], "metadata": { "language_info": { "name": "python" } }, "nbformat": 4, "nbformat_minor": 2 }