{ "cells": [ { "cell_type": "markdown", "metadata": {}, "source": [ "# Advanced Deduplication Strategies in DuckDB & Polars Interoperability Patterns\n", "\n", "This notebook accompanies the blog post on **Advanced Deduplication Strategies in DuckDB & Polars Interoperability Patterns**.\n", "It demonstrates DuckDB's `QUALIFY` clause, `ARG_MAX` aggregations, zero-copy Arrow memory sharing with Polars DataFrames, sliding window debouncing, and out-of-core execution." ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 1. Imports and Zero-Copy Polars Setup" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "import duckdb\n", "import polars as pl\n", "\n", "# Sample DataFrame in Polars\n", "df_pl = pl.DataFrame({\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", "print(df_pl)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 2. DuckDB QUALIFY ROW_NUMBER() Window Deduplication\n", "Query Polars DataFrames directly in SQL using DuckDB's `QUALIFY` clause." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "deduped_rel = duckdb.sql(\"\"\"\n", " SELECT user_id, timestamp, status, amount\n", " FROM df_pl\n", " QUALIFY ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY timestamp DESC) = 1\n", " ORDER BY user_id\n", "\"\"\")\n", "\n", "# Convert DuckDB relation back to Polars DataFrame (zero copy)\n", "df_latest = deduped_rel.pl()\n", "print(df_latest)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 3. DuckDB ARG_MAX Aggregation Deduplication\n", "Extract values at maximum timestamps in a single $O(N)$ pass without window sorting." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "argmax_res = duckdb.sql(\"\"\"\n", " SELECT \n", " user_id,\n", " MAX(timestamp) AS latest_timestamp,\n", " ARG_MAX(status, timestamp) AS latest_status,\n", " ARG_MAX(amount, timestamp) AS latest_amount\n", " FROM df_pl\n", " GROUP BY user_id\n", " ORDER BY user_id\n", "\"\"\").pl()\n", "\n", "print(argmax_res)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 4. Conflict Resolution & Custom Coalesce\n", "Resolve non-key column contradictions across duplicate rows." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "user_updates_pl = 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", "})\n", "\n", "conflict_resolved = duckdb.sql(\"\"\"\n", " SELECT \n", " user_id,\n", " FIRST(email ORDER BY email IS NULL, email) AS email,\n", " FIRST(phone ORDER BY phone IS NULL, phone) AS phone\n", " FROM user_updates_pl\n", " GROUP BY user_id\n", " ORDER BY user_id\n", "\"\"\").pl()\n", "\n", "print(conflict_resolved)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 5. Sliding Window Event Debouncing in DuckDB\n", "Filter rapid retry events occurring within 5 seconds using `LAG()`." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "events_pl = pl.DataFrame({\n", " \"user_id\": [1, 1, 1, 2, 2],\n", " \"timestamp\": [100, 102, 120, 200, 204],\n", " \"event\": [\"click\", \"click\", \"click\", \"buy\", \"buy\"]\n", "})\n", "\n", "debounced_events = duckdb.sql(\"\"\"\n", " WITH ranked AS (\n", " SELECT *,\n", " LAG(timestamp) OVER (PARTITION BY user_id ORDER BY timestamp) AS prev_ts\n", " FROM events_pl\n", " )\n", " SELECT user_id, timestamp, event\n", " FROM ranked\n", " WHERE prev_ts IS NULL OR (timestamp - prev_ts) > 5\n", " ORDER BY user_id, timestamp\n", "\"\"\").pl()\n", "\n", "print(debounced_events)" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "## 6. Out-of-Core Spill-to-Disk Configuration\n", "Set memory limits and temporary spill directories for larger-than-RAM deduplication." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "con = duckdb.connect(database=\":memory:\")\n", "con.execute(\"SET max_memory = '2GB';\")\n", "print(\"Memory limit configured to 2GB for out-of-core spilling.\")" ] } ], "metadata": { "language_info": { "name": "python" } }, "nbformat": 4, "nbformat_minor": 2 }