| |
| <!DOCTYPE html> |
| |
| |
| <html lang="en" data-content_root="../../" data-theme="auto"> |
| |
| <head> |
| <meta charset="utf-8" /> |
| <meta name="viewport" content="width=device-width, initial-scale=1.0" /><meta name="viewport" content="width=device-width, initial-scale=1" /> |
| |
| <title>Aggregation — Apache DataFusion in Python documentation</title> |
| |
| |
| |
| <script data-cfasync="false"> |
| document.documentElement.dataset.mode = localStorage.getItem("mode") || "auto"; |
| document.documentElement.dataset.theme = localStorage.getItem("theme") || "auto"; |
| </script> |
| <!-- |
| this give us a css class that will be invisible only if js is disabled |
| --> |
| <noscript> |
| <style> |
| .pst-js-only { display: none !important; } |
| |
| </style> |
| </noscript> |
| |
| <!-- Loaded before other Sphinx assets --> |
| <link href="../../_static/styles/theme.css?digest=8878045cc6db502f8baf" rel="stylesheet" /> |
| <link href="../../_static/styles/pydata-sphinx-theme.css?digest=8878045cc6db502f8baf" rel="stylesheet" /> |
| |
| <link rel="stylesheet" type="text/css" href="../../_static/pygments.css?v=8f2a1f02" /> |
| <link rel="stylesheet" type="text/css" href="../../_static/mystnb.11b39860a7a0cbfd473a3ad8a317855267ff0bd372690045ca344a6b62be495e.css" /> |
| <link rel="stylesheet" type="text/css" href="../../_static/graphviz.css?v=4ae1632d" /> |
| <link rel="stylesheet" type="text/css" href="../../_static/theme_overrides.css?v=4af573bc" /> |
| |
| <!-- So that users can add custom icons --> |
| <script src="../../_static/scripts/fontawesome.js?digest=8878045cc6db502f8baf"></script> |
| <!-- Pre-loaded scripts that we'll load fully later --> |
| <link rel="preload" as="script" href="../../_static/scripts/bootstrap.js?digest=8878045cc6db502f8baf" /> |
| <link rel="preload" as="script" href="../../_static/scripts/pydata-sphinx-theme.js?digest=8878045cc6db502f8baf" /> |
| |
| <script src="../../_static/documentation_options.js?v=5929fcd5"></script> |
| <script src="../../_static/doctools.js?v=9bcbadda"></script> |
| <script src="../../_static/sphinx_highlight.js?v=dc90522c"></script> |
| <script>DOCUMENTATION_OPTIONS.pagename = 'user-guide/common-operations/aggregations';</script> |
| <script src="../../_static/toc-toggle.js?v=36375f01"></script> |
| <link rel="icon" href="../../_static/favicon.svg"/> |
| <link rel="index" title="Index" href="../../genindex.html" /> |
| <link rel="search" title="Search" href="../../search.html" /> |
| <link rel="next" title="Window Functions" href="windows.html" /> |
| <link rel="prev" title="Spark-Compatible Functions" href="spark-functions.html" /> |
| <meta name="viewport" content="width=device-width, initial-scale=1"/> |
| <meta name="docsearch:language" content="en"/> |
| <meta name="docsearch:version" content="" /> |
| </head> |
| |
| |
| <body data-bs-spy="scroll" data-bs-target=".bd-toc-nav" data-offset="180" data-bs-root-margin="0px 0px -60%" data-default-mode="auto"> |
| |
| |
| |
| <div id="pst-skip-link" class="skip-link d-print-none"><a href="#main-content">Skip to main content</a></div> |
| |
| <div id="pst-scroll-pixel-helper"></div> |
| |
| <button type="button" class="btn rounded-pill" id="pst-back-to-top"> |
| <i class="fa-solid fa-arrow-up"></i>Back to top</button> |
| |
| |
| <dialog id="pst-search-dialog"> |
| |
| <form class="bd-search d-flex align-items-center" |
| action="../../search.html" |
| method="get"> |
| <i class="fa-solid fa-magnifying-glass"></i> |
| <input type="search" |
| class="form-control" |
| name="q" |
| placeholder="Search the docs ..." |
| aria-label="Search the docs ..." |
| autocomplete="off" |
| autocorrect="off" |
| autocapitalize="off" |
| spellcheck="false"/> |
| <span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd>K</kbd></span> |
| </form> |
| </dialog> |
| |
| <div class="pst-async-banner-revealer d-none"> |
| <aside id="bd-header-version-warning" class="d-none d-print-none" aria-label="Version warning"></aside> |
| </div> |
| |
| |
| <header class="bd-header navbar navbar-expand-lg bd-navbar d-print-none"> |
| <div class="bd-header__inner bd-page-width"> |
| <button class="pst-navbar-icon sidebar-toggle primary-toggle" aria-label="Site navigation"> |
| <span class="fa-solid fa-bars"></span> |
| </button> |
| |
| |
| <div class="col-lg-3 navbar-header-items__start"> |
| |
| <div class="navbar-item"> |
| |
| |
| |
| |
| |
| <a class="navbar-brand logo" href="../../index.html"> |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| <img src="../../_static/original.svg" class="logo__image only-light" alt="Apache DataFusion in Python"/> |
| <img src="../../_static/original_dark.svg" class="logo__image only-dark pst-js-only" alt="Apache DataFusion in Python"/> |
| |
| |
| </a></div> |
| |
| </div> |
| |
| <div class="col-lg-9 navbar-header-items"> |
| |
| <div class="me-auto navbar-header-items__center"> |
| |
| <div class="navbar-item"> |
| <nav> |
| <ul class="bd-navbar-elements navbar-nav"> |
| |
| <li class="nav-item current active"> |
| <a class="nav-link nav-internal" href="../index.html"> |
| User Guide |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../contributor-guide/index.html"> |
| Contributor Guide |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../autoapi/index.html"> |
| API Reference |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../links.html"> |
| Links |
| </a> |
| </li> |
| |
| </ul> |
| </nav></div> |
| |
| </div> |
| |
| |
| <div class="navbar-header-items__end"> |
| |
| <div class="navbar-item navbar-persistent--container"> |
| |
| |
| <button class="btn search-button-field search-button__button pst-js-only" title="Search" aria-label="Search" data-bs-placement="bottom" data-bs-toggle="tooltip"> |
| <i class="fa-solid fa-magnifying-glass"></i> |
| <span class="search-button__default-text">Search</span> |
| <span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd class="kbd-shortcut__modifier">K</kbd></span> |
| </button> |
| </div> |
| |
| |
| <div class="navbar-item"><ul class="navbar-icon-links" |
| aria-label="Icon Links"> |
| <li class="nav-item"> |
| |
| |
| |
| |
| |
| |
| |
| |
| <a href="https://github.com/apache/datafusion-python" title="GitHub" class="nav-link pst-navbar-icon" rel="noopener" target="_blank" data-bs-toggle="tooltip" data-bs-placement="bottom"><i class="fa-brands fa-github fa-lg" aria-hidden="true"></i> |
| <span class="sr-only">GitHub</span></a> |
| </li> |
| <li class="nav-item"> |
| |
| |
| |
| |
| |
| |
| |
| |
| <a href="https://docs.rs/datafusion/latest/datafusion/" title="Rust API docs (docs.rs)" class="nav-link pst-navbar-icon" rel="noopener" target="_blank" data-bs-toggle="tooltip" data-bs-placement="bottom"><i class="fa-brands fa-rust fa-lg" aria-hidden="true"></i> |
| <span class="sr-only">Rust API docs (docs.rs)</span></a> |
| </li> |
| </ul></div> |
| |
| <div class="navbar-item"> |
| |
| <button class="btn btn-sm nav-link pst-navbar-icon theme-switch-button pst-js-only" aria-label="Color mode" data-bs-title="Color mode" data-bs-placement="bottom" data-bs-toggle="tooltip"> |
| <i class="theme-switch fa-solid fa-sun fa-lg" data-mode="light" title="Light"></i> |
| <i class="theme-switch fa-solid fa-moon fa-lg" data-mode="dark" title="Dark"></i> |
| <i class="theme-switch fa-solid fa-circle-half-stroke fa-lg" data-mode="auto" title="System Settings"></i> |
| </button></div> |
| |
| </div> |
| |
| </div> |
| |
| |
| <div class="navbar-persistent--mobile"> |
| |
| <button class="btn search-button-field search-button__button pst-js-only" title="Search" aria-label="Search" data-bs-placement="bottom" data-bs-toggle="tooltip"> |
| <i class="fa-solid fa-magnifying-glass"></i> |
| <span class="search-button__default-text">Search</span> |
| <span class="search-button__kbd-shortcut"><kbd class="kbd-shortcut__modifier">Ctrl</kbd>+<kbd class="kbd-shortcut__modifier">K</kbd></span> |
| </button> |
| </div> |
| |
| |
| |
| <button class="pst-navbar-icon sidebar-toggle secondary-toggle" aria-label="On this page"> |
| <span class="fa-solid fa-outdent"></span> |
| </button> |
| |
| </div> |
| |
| </header> |
| |
| |
| <div class="bd-container"> |
| <div class="bd-container__inner bd-page-width"> |
| |
| |
| |
| <dialog id="pst-primary-sidebar-modal"></dialog> |
| <div id="pst-primary-sidebar" class="bd-sidebar-primary bd-sidebar"> |
| |
| |
| |
| <div class="sidebar-header-items sidebar-primary__section"> |
| |
| |
| <div class="sidebar-header-items__center"> |
| |
| |
| |
| <div class="navbar-item"> |
| <nav> |
| <ul class="bd-navbar-elements navbar-nav"> |
| |
| <li class="nav-item current active"> |
| <a class="nav-link nav-internal" href="../index.html"> |
| User Guide |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../contributor-guide/index.html"> |
| Contributor Guide |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../autoapi/index.html"> |
| API Reference |
| </a> |
| </li> |
| |
| |
| <li class="nav-item "> |
| <a class="nav-link nav-internal" href="../../links.html"> |
| Links |
| </a> |
| </li> |
| |
| </ul> |
| </nav></div> |
| |
| |
| </div> |
| |
| |
| |
| <div class="sidebar-header-items__end"> |
| |
| <div class="navbar-item"><ul class="navbar-icon-links" |
| aria-label="Icon Links"> |
| <li class="nav-item"> |
| |
| |
| |
| |
| |
| |
| |
| |
| <a href="https://github.com/apache/datafusion-python" title="GitHub" class="nav-link pst-navbar-icon" rel="noopener" target="_blank" data-bs-toggle="tooltip" data-bs-placement="bottom"><i class="fa-brands fa-github fa-lg" aria-hidden="true"></i> |
| <span class="sr-only">GitHub</span></a> |
| </li> |
| <li class="nav-item"> |
| |
| |
| |
| |
| |
| |
| |
| |
| <a href="https://docs.rs/datafusion/latest/datafusion/" title="Rust API docs (docs.rs)" class="nav-link pst-navbar-icon" rel="noopener" target="_blank" data-bs-toggle="tooltip" data-bs-placement="bottom"><i class="fa-brands fa-rust fa-lg" aria-hidden="true"></i> |
| <span class="sr-only">Rust API docs (docs.rs)</span></a> |
| </li> |
| </ul></div> |
| |
| <div class="navbar-item"> |
| |
| <button class="btn btn-sm nav-link pst-navbar-icon theme-switch-button pst-js-only" aria-label="Color mode" data-bs-title="Color mode" data-bs-placement="bottom" data-bs-toggle="tooltip"> |
| <i class="theme-switch fa-solid fa-sun fa-lg" data-mode="light" title="Light"></i> |
| <i class="theme-switch fa-solid fa-moon fa-lg" data-mode="dark" title="Dark"></i> |
| <i class="theme-switch fa-solid fa-circle-half-stroke fa-lg" data-mode="auto" title="System Settings"></i> |
| </button></div> |
| |
| </div> |
| |
| </div> |
| |
| <div class="sidebar-primary-items__start sidebar-primary__section"> |
| <div class="sidebar-primary-item"> |
| <nav class="bd-docs-nav bd-links" aria-label="Section Navigation"> |
| <p class="bd-links__title" role="heading" aria-level="1">Section Navigation</p> |
| <div class="bd-toc-item navbar-nav"> |
| <ul class="current nav bd-sidenav"> |
| <li class="toctree-l1 current active has-children"><a class="reference internal" href="../index.html">User Guide</a><details open="open"><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul class="current"> |
| <li class="toctree-l2"><a class="reference internal" href="../introduction.html">Introduction</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../basics.html">Concepts</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../data-sources.html">Data Sources</a></li> |
| |
| <li class="toctree-l2 has-children"><a class="reference internal" href="../dataframe/index.html">DataFrames</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l3"><a class="reference internal" href="../dataframe/rendering.html">DataFrame Rendering</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../dataframe/execution-metrics.html">Execution Metrics</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l2 current active has-children"><a class="reference internal" href="index.html">Common Operations</a><details open="open"><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul class="current"> |
| <li class="toctree-l3"><a class="reference internal" href="views.html">Registering Views</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="basic-info.html">Basic Operations</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="select-and-filter.html">Column Selections</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="expressions.html">Expressions</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="joins.html">Joins</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="functions.html">Functions</a></li> |
| |
| <li class="toctree-l3"><a class="reference internal" href="spark-functions.html">Spark-Compatible Functions</a></li> |
| <li class="toctree-l3 current active"><a class="current reference internal" href="#">Aggregation</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="windows.html">Window Functions</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="udf-and-udfa.html">User-Defined Functions</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l2 has-children"><a class="reference internal" href="../io/index.html">IO</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l3"><a class="reference internal" href="../io/arrow.html">Arrow</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../io/avro.html">Avro</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../io/csv.html">CSV</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../io/json.html">JSON</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../io/parquet.html">Parquet</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../io/table_provider.html">Custom Table Provider</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l2"><a class="reference internal" href="../configuration.html">Configuration</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../distributing-work.html">Distributing work</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../sql.html">SQL</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../upgrade-guides.html">Upgrade Guides</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../ai-coding-assistants.html">Using AI Coding Assistants</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l1 has-children"><a class="reference internal" href="../../contributor-guide/index.html">Contributor Guide</a><details open="open"><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l2"><a class="reference internal" href="../../contributor-guide/introduction.html">Introduction</a></li> |
| <li class="toctree-l2"><a class="reference internal" href="../../contributor-guide/ffi.html">Python Extensions</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l1 has-children"><a class="reference internal" href="../../autoapi/index.html">API Reference</a><details open="open"><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l2 has-children"><a class="reference internal" href="../../autoapi/datafusion/index.html">datafusion</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/catalog/index.html">datafusion.catalog</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/context/index.html">datafusion.context</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/dataframe/index.html">datafusion.dataframe</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/dataframe_formatter/index.html">datafusion.dataframe_formatter</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/expr/index.html">datafusion.expr</a></li> |
| <li class="toctree-l3 has-children"><a class="reference internal" href="../../autoapi/datafusion/functions/index.html">datafusion.functions</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l4"><a class="reference internal" href="../../autoapi/datafusion/functions/spark/index.html">datafusion.functions.spark</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l3 has-children"><a class="reference internal" href="../../autoapi/datafusion/input/index.html">datafusion.input</a><details><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l4"><a class="reference internal" href="../../autoapi/datafusion/input/base/index.html">datafusion.input.base</a></li> |
| <li class="toctree-l4"><a class="reference internal" href="../../autoapi/datafusion/input/location/index.html">datafusion.input.location</a></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/io/index.html">datafusion.io</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/ipc/index.html">datafusion.ipc</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/object_store/index.html">datafusion.object_store</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/options/index.html">datafusion.options</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/plan/index.html">datafusion.plan</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/record_batch/index.html">datafusion.record_batch</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/substrait/index.html">datafusion.substrait</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/unparser/index.html">datafusion.unparser</a></li> |
| <li class="toctree-l3"><a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html">datafusion.user_defined</a></li> |
| </ul> |
| </details></li> |
| </ul> |
| </details></li> |
| <li class="toctree-l1 has-children"><a class="reference internal" href="../../links.html">Links</a><details open="open"><summary><span class="toctree-toggle" role="presentation"><i class="fa-solid fa-chevron-down"></i></span></summary><ul> |
| <li class="toctree-l2"><a class="reference external" href="https://github.com/apache/datafusion-python">GitHub and Issue Tracker</a></li> |
| <li class="toctree-l2"><a class="reference external" href="https://docs.rs/datafusion/latest/datafusion/">Rust API Docs</a></li> |
| <li class="toctree-l2"><a class="reference external" href="https://github.com/apache/datafusion/blob/main/CODE_OF_CONDUCT.md">Code of Conduct</a></li> |
| <li class="toctree-l2"><a class="reference external" href="https://github.com/apache/datafusion-python/tree/main/examples">Examples</a></li> |
| </ul> |
| </details></li> |
| </ul> |
| </div> |
| </nav></div> |
| </div> |
| |
| |
| <div class="sidebar-primary-items__end sidebar-primary__section"> |
| <div class="sidebar-primary-item"> |
| <div id="ethical-ad-placement" |
| class="flat" |
| data-ea-publisher="readthedocs" |
| data-ea-type="readthedocs-sidebar" |
| data-ea-manual="true"> |
| </div></div> |
| </div> |
| |
| |
| </div> |
| |
| <main id="main-content" class="bd-main" role="main"> |
| |
| |
| <div class="bd-content"> |
| <div class="bd-article-container"> |
| |
| <div class="bd-header-article d-print-none"> |
| <div class="header-article-items header-article__inner"> |
| |
| <div class="header-article-items__start"> |
| |
| <div class="header-article-item"> |
| |
| <nav aria-label="Breadcrumb" class="d-print-none"> |
| <ul class="bd-breadcrumbs"> |
| |
| <li class="breadcrumb-item breadcrumb-home"> |
| <a href="../../index.html" class="nav-link" aria-label="Home"> |
| <i class="fa-solid fa-home"></i> |
| </a> |
| </li> |
| |
| <li class="breadcrumb-item"><a href="../index.html" class="nav-link">User Guide</a></li> |
| |
| |
| <li class="breadcrumb-item"><a href="index.html" class="nav-link">Common Operations</a></li> |
| |
| <li class="breadcrumb-item active" aria-current="page"><span class="ellipsis">Aggregation</span></li> |
| </ul> |
| </nav> |
| </div> |
| |
| </div> |
| |
| |
| </div> |
| </div> |
| |
| |
| |
| |
| <div id="searchbox"></div> |
| <article class="bd-article"> |
| |
| <!--- |
| Licensed to the Apache Software Foundation (ASF) under one |
| or more contributor license agreements. See the NOTICE file |
| distributed with this work for additional information |
| regarding copyright ownership. The ASF licenses this file |
| to you 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. |
| --> |
| <section id="aggregation"> |
| <span id="id1"></span><h1>Aggregation<a class="headerlink" href="#aggregation" title="Link to this heading">#</a></h1> |
| <p>An aggregate or aggregation is a function where the values of multiple rows are processed together |
| to form a single summary value. For performing an aggregation, DataFusion provides the |
| <a class="reference internal" href="../../autoapi/datafusion/dataframe/index.html#datafusion.dataframe.DataFrame.aggregate" title="datafusion.dataframe.DataFrame.aggregate"><code class="xref py py-func docutils literal notranslate"><span class="pre">aggregate()</span></code></a></p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="kn">from</span><span class="w"> </span><span class="nn">datafusion</span><span class="w"> </span><span class="kn">import</span> <span class="n">SessionContext</span><span class="p">,</span> <span class="n">col</span><span class="p">,</span> <span class="n">lit</span><span class="p">,</span> <span class="n">functions</span> <span class="k">as</span> <span class="n">f</span> |
| |
| <span class="n">ctx</span> <span class="o">=</span> <span class="n">SessionContext</span><span class="p">()</span> |
| <span class="n">df</span> <span class="o">=</span> <span class="n">ctx</span><span class="o">.</span><span class="n">read_csv</span><span class="p">(</span><span class="s2">"pokemon.csv"</span><span class="p">)</span> |
| |
| <span class="n">col_type_1</span> <span class="o">=</span> <span class="n">col</span><span class="p">(</span><span class="s1">'"Type 1"'</span><span class="p">)</span> |
| <span class="n">col_type_2</span> <span class="o">=</span> <span class="n">col</span><span class="p">(</span><span class="s1">'"Type 2"'</span><span class="p">)</span> |
| <span class="n">col_speed</span> <span class="o">=</span> <span class="n">col</span><span class="p">(</span><span class="s1">'"Speed"'</span><span class="p">)</span> |
| <span class="n">col_attack</span> <span class="o">=</span> <span class="n">col</span><span class="p">(</span><span class="s1">'"Attack"'</span><span class="p">)</span> |
| |
| <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">approx_distinct</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">approx_median</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Median Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">approx_percentile_cont</span><span class="p">(</span><span class="n">col_speed</span><span class="p">,</span> <span class="mf">0.9</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"90% Speed"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+-------+--------------+--------------------+ |
| | Type 1 | Count | Median Speed | 90% Speed | |
| +----------+-------+--------------+--------------------+ |
| | Water | 21 | 70.0 | 90.0 | |
| | Rock | 8 | 55.0 | 140.0 | |
| | Ghost | 4 | 101.25 | 130.0 | |
| | Ice | 2 | 90.0 | 95.0 | |
| | Dragon | 3 | 70.0 | 80.0 | |
| | Grass | 8 | 55.0 | 80.0 | |
| | Fire | 8 | 91.75 | 100.25 | |
| | Normal | 20 | 71.0 | 110.70000000000002 | |
| | Poison | 12 | 55.0 | 85.5 | |
| | Fighting | 7 | 70.0 | 93.4 | |
| +----------+-------+--------------+--------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>When <code class="code docutils literal notranslate"><span class="pre">group_by</span></code> is <code class="code docutils literal notranslate"><span class="pre">None</span></code> or an empty list, the aggregation is done over the whole |
| <a class="reference internal" href="../../autoapi/datafusion/dataframe/index.html#datafusion.dataframe.DataFrame" title="datafusion.dataframe.DataFrame"><code class="xref py py-class docutils literal notranslate"><span class="pre">DataFrame</span></code></a>. For grouping the <code class="code docutils literal notranslate"><span class="pre">group_by</span></code> list must contain at least one column.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">max</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Max Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">min</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Min Speed"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+-----------+--------------------+-----------+ |
| | Type 1 | Max Speed | Avg Speed | Min Speed | |
| +----------+-----------+--------------------+-----------+ |
| | Water | 115 | 67.25806451612904 | 15 | |
| | Rock | 150 | 67.5 | 20 | |
| | Ghost | 130 | 103.75 | 80 | |
| | Ice | 95 | 90.0 | 85 | |
| | Dragon | 80 | 66.66666666666667 | 50 | |
| | Grass | 80 | 54.23076923076923 | 30 | |
| | Fire | 105 | 86.28571428571429 | 60 | |
| | Normal | 121 | 72.75 | 20 | |
| | Poison | 90 | 58.785714285714285 | 25 | |
| | Fighting | 95 | 66.14285714285714 | 35 | |
| +----------+-----------+--------------------+-----------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>More than one column can be used for grouping</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">,</span> <span class="n">col_type_2</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">max</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Max Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">min</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Min Speed"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +--------+---------+-----------+-------------------+-----------+ |
| | Type 1 | Type 2 | Max Speed | Avg Speed | Min Speed | |
| +--------+---------+-----------+-------------------+-----------+ |
| | Water | | 90 | 68.05263157894737 | 40 | |
| | Poison | Ground | 85 | 80.5 | 76 | |
| | Grass | Psychic | 55 | 47.5 | 40 | |
| | Water | Flying | 81 | 81.0 | 81 | |
| | Rock | Flying | 150 | 140.0 | 130 | |
| | Ice | Flying | 85 | 85.0 | 85 | |
| | Dragon | | 70 | 60.0 | 50 | |
| | Dragon | Flying | 80 | 80.0 | 80 | |
| | Fire | | 105 | 81.8 | 60 | |
| | Fire | Flying | 100 | 96.66666666666667 | 90 | |
| +--------+---------+-----------+-------------------+-----------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <section id="setting-parameters"> |
| <h2>Setting Parameters<a class="headerlink" href="#setting-parameters" title="Link to this heading">#</a></h2> |
| <p>Each of the built in aggregate functions provides arguments for the parameters that affect their |
| operation. These can also be overridden using the builder approach to setting any of the following |
| parameters. When you use the builder, you must call <code class="docutils literal notranslate"><span class="pre">build()</span></code> to finish. For example, these two |
| expressions are equivalent.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">first_1</span> <span class="o">=</span> <span class="n">f</span><span class="o">.</span><span class="n">first_value</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"a"</span><span class="p">),</span> <span class="n">order_by</span><span class="o">=</span><span class="p">[</span><span class="n">col</span><span class="p">(</span><span class="s2">"a"</span><span class="p">)])</span> |
| <span class="n">first_2</span> <span class="o">=</span> <span class="n">f</span><span class="o">.</span><span class="n">first_value</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"a"</span><span class="p">))</span><span class="o">.</span><span class="n">order_by</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"a"</span><span class="p">))</span><span class="o">.</span><span class="n">build</span><span class="p">()</span> |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <section id="ordering"> |
| <h3>Ordering<a class="headerlink" href="#ordering" title="Link to this heading">#</a></h3> |
| <p>You can control the order in which rows are processed by window functions by providing |
| a list of <code class="docutils literal notranslate"><span class="pre">order_by</span></code> functions for the <code class="docutils literal notranslate"><span class="pre">order_by</span></code> parameter. In the following example, we |
| sort the Pokemon by their attack in increasing order and take the first value, which gives us the |
| Pokemon with the smallest attack value in each <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code>.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">col</span><span class="p">(</span><span class="s1">'"Type 1"'</span><span class="p">)],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">first_value</span><span class="p">(</span> |
| <span class="n">col</span><span class="p">(</span><span class="s1">'"Name"'</span><span class="p">),</span> |
| <span class="n">order_by</span><span class="o">=</span><span class="p">[</span><span class="n">col</span><span class="p">(</span><span class="s1">'"Attack"'</span><span class="p">)</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">)]</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Smallest Attack"</span><span class="p">)</span> |
| <span class="p">])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+-----------------+ |
| | Type 1 | Smallest Attack | |
| +----------+-----------------+ |
| | Water | Magikarp | |
| | Rock | Omanyte | |
| | Ghost | Gastly | |
| | Ice | Jynx | |
| | Dragon | Dratini | |
| | Grass | Exeggcute | |
| | Fire | Vulpix | |
| | Normal | Chansey | |
| | Poison | Zubat | |
| | Fighting | Mankey | |
| +----------+-----------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| </section> |
| <section id="distinct"> |
| <h3>Distinct<a class="headerlink" href="#distinct" title="Link to this heading">#</a></h3> |
| <p>When you set the parameter <code class="docutils literal notranslate"><span class="pre">distinct</span></code> to <code class="docutils literal notranslate"><span class="pre">True</span></code>, then unique values will only be evaluated one |
| time each. Suppose we want to create an array of all of the <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> for each <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> of our |
| Pokemon set. Since there will be many entries of <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> we only one each distinct value.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">array_agg</span><span class="p">(</span><span class="n">col_type_2</span><span class="p">,</span> <span class="n">distinct</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Type 2 List"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+--------------------------------------------------+ |
| | Type 1 | Type 2 List | |
| +----------+--------------------------------------------------+ |
| | Water | [Fighting, Flying, , Poison, Psychic, Dark, Ice] | |
| | Rock | [Water, Ground, Flying] | |
| | Ghost | [Poison] | |
| | Ice | [Flying, Psychic] | |
| | Dragon | [, Flying] | |
| | Grass | [Psychic, , Poison] | |
| | Fire | [, Dragon, Flying] | |
| | Normal | [Fairy, Flying, ] | |
| | Poison | [Ground, Flying, ] | |
| | Fighting | [] | |
| +----------+--------------------------------------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>In the output of the above we can see that there are some <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> for which the <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> entry |
| is <code class="docutils literal notranslate"><span class="pre">null</span></code>. In reality, we probably want to filter those out. We can do this in two ways. First, |
| we can filter DataFrame rows that have no <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code>. If we do this, we might have some <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> |
| entries entirely removed. The second is we can use the <code class="docutils literal notranslate"><span class="pre">filter</span></code> argument described below.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">filter</span><span class="p">(</span><span class="n">col_type_2</span><span class="o">.</span><span class="n">is_not_null</span><span class="p">())</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">array_agg</span><span class="p">(</span><span class="n">col_type_2</span><span class="p">,</span> <span class="n">distinct</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Type 2 List"</span><span class="p">)])</span> |
| |
| <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">array_agg</span><span class="p">(</span><span class="n">col_type_2</span><span class="p">,</span> <span class="n">distinct</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="nb">filter</span><span class="o">=</span><span class="n">col_type_2</span><span class="o">.</span><span class="n">is_not_null</span><span class="p">())</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Type 2 List"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+------------------------------------------------+ |
| | Type 1 | Type 2 List | |
| +----------+------------------------------------------------+ |
| | Water | [Fighting, Ice, Flying, Psychic, Dark, Poison] | |
| | Rock | [Flying, Ground, Water] | |
| | Ghost | [Poison] | |
| | Ice | [Psychic, Flying] | |
| | Dragon | [Flying] | |
| | Grass | [Psychic, Poison] | |
| | Fire | [Flying, Dragon] | |
| | Normal | [Fairy, Flying] | |
| | Poison | [Flying, Ground] | |
| | Fighting | | |
| +----------+------------------------------------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>Which approach you take should depend on your use case.</p> |
| </section> |
| <section id="null-treatment"> |
| <h3>Null Treatment<a class="headerlink" href="#null-treatment" title="Link to this heading">#</a></h3> |
| <p>This option allows you to either respect or ignore null values.</p> |
| <p>One common usage for handling nulls is the case where you want to find the first value within a |
| partition. By setting the null treatment to ignore nulls, we can find the first non-null value |
| in our partition.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="kn">from</span><span class="w"> </span><span class="nn">datafusion.common</span><span class="w"> </span><span class="kn">import</span> <span class="n">NullTreatment</span> |
| |
| <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">first_value</span><span class="p">(</span> |
| <span class="n">col_type_2</span><span class="p">,</span> |
| <span class="n">order_by</span><span class="o">=</span><span class="p">[</span><span class="n">col_attack</span><span class="p">],</span> |
| <span class="n">null_treatment</span><span class="o">=</span><span class="n">NullTreatment</span><span class="o">.</span><span class="n">RESPECT_NULLS</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Lowest Attack Type 2"</span><span class="p">)])</span> |
| |
| <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">first_value</span><span class="p">(</span> |
| <span class="n">col_type_2</span><span class="p">,</span> |
| <span class="n">order_by</span><span class="o">=</span><span class="p">[</span><span class="n">col_attack</span><span class="p">],</span> |
| <span class="n">null_treatment</span><span class="o">=</span><span class="n">NullTreatment</span><span class="o">.</span><span class="n">IGNORE_NULLS</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Lowest Attack Type 2"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+----------------------+ |
| | Type 1 | Lowest Attack Type 2 | |
| +----------+----------------------+ |
| | Water | Poison | |
| | Rock | Water | |
| | Ghost | Poison | |
| | Ice | Psychic | |
| | Dragon | Flying | |
| | Grass | Psychic | |
| | Fire | Flying | |
| | Normal | Flying | |
| | Poison | Flying | |
| | Fighting | | |
| +----------+----------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| </section> |
| <section id="filter"> |
| <h3>Filter<a class="headerlink" href="#filter" title="Link to this heading">#</a></h3> |
| <p>Using the filter option is useful for filtering results to include in the aggregate function. It can |
| be seen in the example above on how this can be useful to only filter rows evaluated by the |
| aggregate function without filtering rows from the entire DataFrame.</p> |
| <p>Filter takes a single expression.</p> |
| <p>Suppose we want to find the speed values for only Pokemon that have low Attack values.</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed All"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">,</span> <span class="nb">filter</span><span class="o">=</span><span class="n">col_attack</span> <span class="o"><</span> <span class="n">lit</span><span class="p">(</span><span class="mi">50</span><span class="p">))</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed Low Attack"</span><span class="p">)])</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+--------------------+----------------------+ |
| | Type 1 | Avg Speed All | Avg Speed Low Attack | |
| +----------+--------------------+----------------------+ |
| | Water | 67.25806451612904 | 63.833333333333336 | |
| | Rock | 67.5 | 52.5 | |
| | Ghost | 103.75 | 80.0 | |
| | Ice | 90.0 | | |
| | Dragon | 66.66666666666667 | | |
| | Grass | 54.23076923076923 | 42.5 | |
| | Fire | 86.28571428571429 | 65.0 | |
| | Normal | 72.75 | 52.8 | |
| | Poison | 58.785714285714285 | 48.0 | |
| | Fighting | 66.14285714285714 | | |
| +----------+--------------------+----------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| </section> |
| <section id="comparing-subsets-within-a-group"> |
| <h3>Comparing subsets within a group<a class="headerlink" href="#comparing-subsets-within-a-group" title="Link to this heading">#</a></h3> |
| <p>Sometimes you need to compare the full membership of a group against a |
| subset that meets some condition — for example, “which groups have at least |
| one failure, but not every member failed?”. The <code class="docutils literal notranslate"><span class="pre">filter</span></code> argument on an |
| aggregate restricts the rows that contribute to <em>that</em> aggregate without |
| dropping the group, so a single pass can produce both the full set and the |
| filtered subset side by side. Pairing |
| <a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.array_agg" title="datafusion.functions.array_agg"><code class="xref py py-func docutils literal notranslate"><span class="pre">array_agg()</span></code></a> with <code class="docutils literal notranslate"><span class="pre">distinct=True</span></code> and |
| <code class="docutils literal notranslate"><span class="pre">filter=</span></code> is a compact way to express this: collect the distinct values |
| of the group, collect the distinct values that satisfy the condition, then |
| compare the two arrays.</p> |
| <p>Suppose each row records a line item with the supplier that fulfilled it and |
| a flag for whether that supplier met the commit date. We want to identify |
| <em>partially failed</em> orders — orders where at least one supplier failed but |
| not every supplier failed:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">orders_df</span> <span class="o">=</span> <span class="n">ctx</span><span class="o">.</span><span class="n">from_pydict</span><span class="p">(</span> |
| <span class="p">{</span> |
| <span class="s2">"order_id"</span><span class="p">:</span> <span class="p">[</span><span class="mi">1</span><span class="p">,</span> <span class="mi">1</span><span class="p">,</span> <span class="mi">1</span><span class="p">,</span> <span class="mi">2</span><span class="p">,</span> <span class="mi">2</span><span class="p">,</span> <span class="mi">3</span><span class="p">,</span> <span class="mi">4</span><span class="p">,</span> <span class="mi">4</span><span class="p">],</span> |
| <span class="s2">"supplier_id"</span><span class="p">:</span> <span class="p">[</span><span class="mi">100</span><span class="p">,</span> <span class="mi">101</span><span class="p">,</span> <span class="mi">102</span><span class="p">,</span> <span class="mi">200</span><span class="p">,</span> <span class="mi">201</span><span class="p">,</span> <span class="mi">300</span><span class="p">,</span> <span class="mi">400</span><span class="p">,</span> <span class="mi">401</span><span class="p">],</span> |
| <span class="s2">"failed"</span><span class="p">:</span> <span class="p">[</span><span class="kc">False</span><span class="p">,</span> <span class="kc">True</span><span class="p">,</span> <span class="kc">False</span><span class="p">,</span> <span class="kc">False</span><span class="p">,</span> <span class="kc">False</span><span class="p">,</span> <span class="kc">True</span><span class="p">,</span> <span class="kc">True</span><span class="p">,</span> <span class="kc">True</span><span class="p">],</span> |
| <span class="p">},</span> |
| <span class="p">)</span> |
| |
| <span class="n">grouped</span> <span class="o">=</span> <span class="n">orders_df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">col</span><span class="p">(</span><span class="s2">"order_id"</span><span class="p">)],</span> |
| <span class="p">[</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">array_agg</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"supplier_id"</span><span class="p">),</span> <span class="n">distinct</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"all_suppliers"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">array_agg</span><span class="p">(</span> |
| <span class="n">col</span><span class="p">(</span><span class="s2">"supplier_id"</span><span class="p">),</span> |
| <span class="nb">filter</span><span class="o">=</span><span class="n">col</span><span class="p">(</span><span class="s2">"failed"</span><span class="p">),</span> |
| <span class="n">distinct</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"failed_suppliers"</span><span class="p">),</span> |
| <span class="p">],</span> |
| <span class="p">)</span> |
| |
| <span class="n">grouped</span><span class="o">.</span><span class="n">filter</span><span class="p">(</span> |
| <span class="p">(</span><span class="n">f</span><span class="o">.</span><span class="n">array_length</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"failed_suppliers"</span><span class="p">))</span> <span class="o">></span> <span class="n">lit</span><span class="p">(</span><span class="mi">0</span><span class="p">))</span> |
| <span class="o">&</span> <span class="p">(</span><span class="n">f</span><span class="o">.</span><span class="n">array_length</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"failed_suppliers"</span><span class="p">))</span> <span class="o"><</span> <span class="n">f</span><span class="o">.</span><span class="n">array_length</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"all_suppliers"</span><span class="p">)))</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">select</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">"order_id"</span><span class="p">),</span> <span class="n">col</span><span class="p">(</span><span class="s2">"failed_suppliers"</span><span class="p">))</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+------------------+ |
| | order_id | failed_suppliers | |
| +----------+------------------+ |
| | 1 | [101] | |
| +----------+------------------+ |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>Order 1 is partial (one of three suppliers failed). Order 2 is excluded |
| because no supplier failed, order 3 because its only supplier failed, and |
| order 4 because both of its suppliers failed.</p> |
| </section> |
| </section> |
| <section id="grouping-sets"> |
| <h2>Grouping Sets<a class="headerlink" href="#grouping-sets" title="Link to this heading">#</a></h2> |
| <p>The default style of aggregation produces one row per group. Sometimes you want a single query to |
| produce rows at multiple levels of detail — for example, totals per type <em>and</em> an overall grand |
| total, or subtotals for every combination of two columns plus the individual column totals. Writing |
| separate queries and concatenating them is tedious and runs the data multiple times. Grouping sets |
| solve this by letting you specify several grouping levels in one pass.</p> |
| <p>DataFusion supports three grouping set styles through the |
| <a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet" title="datafusion.expr.GroupingSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">GroupingSet</span></code></a> class:</p> |
| <ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.rollup" title="datafusion.expr.GroupingSet.rollup"><code class="xref py py-meth docutils literal notranslate"><span class="pre">rollup()</span></code></a> — hierarchical subtotals, like a drill-down report</p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.cube" title="datafusion.expr.GroupingSet.cube"><code class="xref py py-meth docutils literal notranslate"><span class="pre">cube()</span></code></a> — every possible subtotal combination, like a pivot table</p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.grouping_sets" title="datafusion.expr.GroupingSet.grouping_sets"><code class="xref py py-meth docutils literal notranslate"><span class="pre">grouping_sets()</span></code></a> — explicitly list exactly which grouping levels you want</p></li> |
| </ul> |
| <p>Because result rows come from different grouping levels, a column that is <em>not</em> part of a |
| particular level will be <code class="docutils literal notranslate"><span class="pre">null</span></code> in that row. Use <a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.grouping" title="datafusion.functions.grouping"><code class="xref py py-func docutils literal notranslate"><span class="pre">grouping()</span></code></a> to |
| distinguish a real <code class="docutils literal notranslate"><span class="pre">null</span></code> in the data from one that means “this column was aggregated across.” |
| It returns <code class="docutils literal notranslate"><span class="pre">0</span></code> when the column is a grouping key for that row, and <code class="docutils literal notranslate"><span class="pre">1</span></code> when it is not.</p> |
| <section id="rollup"> |
| <h3>Rollup<a class="headerlink" href="#rollup" title="Link to this heading">#</a></h3> |
| <p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.rollup" title="datafusion.expr.GroupingSet.rollup"><code class="xref py py-meth docutils literal notranslate"><span class="pre">rollup()</span></code></a> creates a hierarchy. <code class="docutils literal notranslate"><span class="pre">rollup(a,</span> <span class="pre">b)</span></code> produces |
| grouping sets <code class="docutils literal notranslate"><span class="pre">(a,</span> <span class="pre">b)</span></code>, <code class="docutils literal notranslate"><span class="pre">(a)</span></code>, and <code class="docutils literal notranslate"><span class="pre">()</span></code> — like nested subtotals in a report. This is useful |
| when your columns have a natural hierarchy, such as region → city or type → subtype.</p> |
| <p>Suppose we want to summarize Pokemon stats by <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> with subtotals and a grand total. With |
| the default aggregation style we would need two separate queries. With <code class="docutils literal notranslate"><span class="pre">rollup</span></code> we get it all at |
| once:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="kn">from</span><span class="w"> </span><span class="nn">datafusion.expr</span><span class="w"> </span><span class="kn">import</span> <span class="n">GroupingSet</span> |
| |
| <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">rollup</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">)],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">max</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Max Speed"</span><span class="p">)]</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">))</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+-------+-------------------+-----------+ |
| | Type 1 | Count | Avg Speed | Max Speed | |
| +----------+-------+-------------------+-----------+ |
| | | 163 | 71.65030674846626 | 150 | |
| | Bug | 14 | 66.78571428571429 | 145 | |
| | Dragon | 3 | 66.66666666666667 | 80 | |
| | Electric | 9 | 98.88888888888889 | 140 | |
| | Fairy | 2 | 47.5 | 60 | |
| | Fighting | 7 | 66.14285714285714 | 95 | |
| | Fire | 14 | 86.28571428571429 | 105 | |
| | Ghost | 4 | 103.75 | 130 | |
| | Grass | 13 | 54.23076923076923 | 80 | |
| | Ground | 8 | 58.125 | 120 | |
| +----------+-------+-------------------+-----------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>The first row — where <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> is <code class="docutils literal notranslate"><span class="pre">null</span></code> — is the grand total across all types. But how do you |
| tell a grand-total <code class="docutils literal notranslate"><span class="pre">null</span></code> apart from a Pokemon that genuinely has no type? The |
| <a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.grouping" title="datafusion.functions.grouping"><code class="xref py py-func docutils literal notranslate"><span class="pre">grouping()</span></code></a> function returns <code class="docutils literal notranslate"><span class="pre">0</span></code> when the column is a grouping key |
| for that row and <code class="docutils literal notranslate"><span class="pre">1</span></code> when it is aggregated across.</p> |
| <p>Apply <code class="docutils literal notranslate"><span class="pre">.alias()</span></code> to the <code class="docutils literal notranslate"><span class="pre">grouping()</span></code> expression to give the column a readable name:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">result</span> <span class="o">=</span> <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">rollup</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">)],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">grouping</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Is Total"</span><span class="p">)]</span> |
| <span class="p">)</span> |
| <span class="n">result</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">))</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+-------+-------------------+----------+ |
| | Type 1 | Count | Avg Speed | Is Total | |
| +----------+-------+-------------------+----------+ |
| | | 163 | 71.65030674846626 | 1 | |
| | Bug | 14 | 66.78571428571429 | 0 | |
| | Dragon | 3 | 66.66666666666667 | 0 | |
| | Electric | 9 | 98.88888888888889 | 0 | |
| | Fairy | 2 | 47.5 | 0 | |
| | Fighting | 7 | 66.14285714285714 | 0 | |
| | Fire | 14 | 86.28571428571429 | 0 | |
| | Ghost | 4 | 103.75 | 0 | |
| | Grass | 13 | 54.23076923076923 | 0 | |
| | Ground | 8 | 58.125 | 0 | |
| +----------+-------+-------------------+----------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>With two columns the hierarchy becomes more apparent. <code class="docutils literal notranslate"><span class="pre">rollup(Type</span> <span class="pre">1,</span> <span class="pre">Type</span> <span class="pre">2)</span></code> produces:</p> |
| <ul class="simple"> |
| <li><p>one row per <code class="docutils literal notranslate"><span class="pre">(Type</span> <span class="pre">1,</span> <span class="pre">Type</span> <span class="pre">2)</span></code> pair — the most detailed level</p></li> |
| <li><p>one row per <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> — subtotals</p></li> |
| <li><p>one grand total row</p></li> |
| </ul> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">rollup</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">,</span> <span class="n">col_type_2</span><span class="p">)],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">)]</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span> |
| <span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">),</span> |
| <span class="n">col_type_2</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span> |
| <span class="p">)</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +----------+--------+-------+--------------------+ |
| | Type 1 | Type 2 | Count | Avg Speed | |
| +----------+--------+-------+--------------------+ |
| | | | 163 | 71.65030674846626 | |
| | Bug | | 3 | 53.333333333333336 | |
| | Bug | | 14 | 66.78571428571429 | |
| | Bug | Flying | 3 | 93.33333333333333 | |
| | Bug | Grass | 2 | 27.5 | |
| | Bug | Poison | 6 | 73.33333333333333 | |
| | Dragon | | 3 | 66.66666666666667 | |
| | Dragon | | 2 | 60.0 | |
| | Dragon | Flying | 1 | 80.0 | |
| | Electric | | 6 | 112.5 | |
| +----------+--------+-------+--------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| </section> |
| <section id="cube"> |
| <h3>Cube<a class="headerlink" href="#cube" title="Link to this heading">#</a></h3> |
| <p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.cube" title="datafusion.expr.GroupingSet.cube"><code class="xref py py-meth docutils literal notranslate"><span class="pre">cube()</span></code></a> produces every possible subset. <code class="docutils literal notranslate"><span class="pre">cube(a,</span> <span class="pre">b)</span></code> |
| produces grouping sets <code class="docutils literal notranslate"><span class="pre">(a,</span> <span class="pre">b)</span></code>, <code class="docutils literal notranslate"><span class="pre">(a)</span></code>, <code class="docutils literal notranslate"><span class="pre">(b)</span></code>, and <code class="docutils literal notranslate"><span class="pre">()</span></code> — one more than <code class="docutils literal notranslate"><span class="pre">rollup</span></code> because |
| it also includes <code class="docutils literal notranslate"><span class="pre">(b)</span></code> alone. This is useful when neither column is “above” the other in a |
| hierarchy and you want all cross-tabulations.</p> |
| <p>For our Pokemon data, <code class="docutils literal notranslate"><span class="pre">cube(Type</span> <span class="pre">1,</span> <span class="pre">Type</span> <span class="pre">2)</span></code> gives us stats broken down by the type pair, |
| by <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> alone, by <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> alone, and a grand total — all in one query:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">cube</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">,</span> <span class="n">col_type_2</span><span class="p">)],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">)]</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span> |
| <span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">),</span> |
| <span class="n">col_type_2</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span> |
| <span class="p">)</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +--------+----------+-------+--------------------+ |
| | Type 1 | Type 2 | Count | Avg Speed | |
| +--------+----------+-------+--------------------+ |
| | | | 86 | 72.46511627906976 | |
| | | | 163 | 71.65030674846626 | |
| | | Dark | 1 | 81.0 | |
| | | Dragon | 1 | 100.0 | |
| | | Fairy | 3 | 51.666666666666664 | |
| | | Fighting | 1 | 70.0 | |
| | | Flying | 23 | 91.08695652173913 | |
| | | Grass | 2 | 27.5 | |
| | | Ground | 6 | 55.166666666666664 | |
| | | Ice | 3 | 66.66666666666667 | |
| +--------+----------+-------+--------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>Compared to the <code class="docutils literal notranslate"><span class="pre">rollup</span></code> example above, notice the extra rows where <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> is <code class="docutils literal notranslate"><span class="pre">null</span></code> but |
| <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> has a value — those are the per-<code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> subtotals that <code class="docutils literal notranslate"><span class="pre">rollup</span></code> does not include.</p> |
| </section> |
| <section id="explicit-grouping-sets"> |
| <h3>Explicit Grouping Sets<a class="headerlink" href="#explicit-grouping-sets" title="Link to this heading">#</a></h3> |
| <p><a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.grouping_sets" title="datafusion.expr.GroupingSet.grouping_sets"><code class="xref py py-meth docutils literal notranslate"><span class="pre">grouping_sets()</span></code></a> lets you list exactly which grouping levels |
| you need when <code class="docutils literal notranslate"><span class="pre">rollup</span></code> or <code class="docutils literal notranslate"><span class="pre">cube</span></code> would produce too many or too few. Each argument is a list of |
| columns forming one grouping set.</p> |
| <p>For example, if we want only the per-<code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> totals and per-<code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> totals — but <em>not</em> the |
| full <code class="docutils literal notranslate"><span class="pre">(Type</span> <span class="pre">1,</span> <span class="pre">Type</span> <span class="pre">2)</span></code> detail rows or the grand total — we can ask for exactly that:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">grouping_sets</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span><span class="n">col_type_2</span><span class="p">])],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">)]</span> |
| <span class="p">)</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span> |
| <span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">),</span> |
| <span class="n">col_type_2</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span> |
| <span class="p">)</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +--------+----------+-------+--------------------+ |
| | Type 1 | Type 2 | Count | Avg Speed | |
| +--------+----------+-------+--------------------+ |
| | | | 86 | 72.46511627906976 | |
| | | Dark | 1 | 81.0 | |
| | | Dragon | 1 | 100.0 | |
| | | Fairy | 3 | 51.666666666666664 | |
| | | Fighting | 1 | 70.0 | |
| | | Flying | 23 | 91.08695652173913 | |
| | | Grass | 2 | 27.5 | |
| | | Ground | 6 | 55.166666666666664 | |
| | | Ice | 3 | 66.66666666666667 | |
| | | Poison | 22 | 71.5909090909091 | |
| +--------+----------+-------+--------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>Each row belongs to exactly one grouping level. The <a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.grouping" title="datafusion.functions.grouping"><code class="xref py py-func docutils literal notranslate"><span class="pre">grouping()</span></code></a> |
| function tells you which level each row comes from:</p> |
| <div class="cell docutils container"> |
| <div class="cell_input docutils container"> |
| <div class="highlight-ipython3 notranslate"><div class="highlight"><pre><span></span><span class="n">result</span> <span class="o">=</span> <span class="n">df</span><span class="o">.</span><span class="n">aggregate</span><span class="p">(</span> |
| <span class="p">[</span><span class="n">GroupingSet</span><span class="o">.</span><span class="n">grouping_sets</span><span class="p">([</span><span class="n">col_type_1</span><span class="p">],</span> <span class="p">[</span><span class="n">col_type_2</span><span class="p">])],</span> |
| <span class="p">[</span><span class="n">f</span><span class="o">.</span><span class="n">count</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Count"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">avg</span><span class="p">(</span><span class="n">col_speed</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"Avg Speed"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">grouping</span><span class="p">(</span><span class="n">col_type_1</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"grouping(Type 1)"</span><span class="p">),</span> |
| <span class="n">f</span><span class="o">.</span><span class="n">grouping</span><span class="p">(</span><span class="n">col_type_2</span><span class="p">)</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">"grouping(Type 2)"</span><span class="p">)]</span> |
| <span class="p">)</span> |
| <span class="n">result</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span> |
| <span class="n">col_type_1</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">),</span> |
| <span class="n">col_type_2</span><span class="o">.</span><span class="n">sort</span><span class="p">(</span><span class="n">ascending</span><span class="o">=</span><span class="kc">True</span><span class="p">,</span> <span class="n">nulls_first</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span> |
| <span class="p">)</span> |
| </pre></div> |
| </div> |
| </div> |
| <div class="cell_output docutils container"> |
| <div class="output text_plain highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame() |
| +--------+----------+-------+--------------------+------------------+------------------+ |
| | Type 1 | Type 2 | Count | Avg Speed | grouping(Type 1) | grouping(Type 2) | |
| +--------+----------+-------+--------------------+------------------+------------------+ |
| | | | 86 | 72.46511627906976 | 1 | 0 | |
| | | Dark | 1 | 81.0 | 1 | 0 | |
| | | Dragon | 1 | 100.0 | 1 | 0 | |
| | | Fairy | 3 | 51.666666666666664 | 1 | 0 | |
| | | Fighting | 1 | 70.0 | 1 | 0 | |
| | | Flying | 23 | 91.08695652173913 | 1 | 0 | |
| | | Grass | 2 | 27.5 | 1 | 0 | |
| | | Ground | 6 | 55.166666666666664 | 1 | 0 | |
| | | Ice | 3 | 66.66666666666667 | 1 | 0 | |
| | | Poison | 22 | 71.5909090909091 | 1 | 0 | |
| +--------+----------+-------+--------------------+------------------+------------------+ |
| Data truncated. |
| </pre></div> |
| </div> |
| </div> |
| </div> |
| <p>Where <code class="docutils literal notranslate"><span class="pre">grouping(Type</span> <span class="pre">1)</span></code> is <code class="docutils literal notranslate"><span class="pre">0</span></code> the row is a per-<code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> total (and <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> is <code class="docutils literal notranslate"><span class="pre">null</span></code>). |
| Where <code class="docutils literal notranslate"><span class="pre">grouping(Type</span> <span class="pre">2)</span></code> is <code class="docutils literal notranslate"><span class="pre">0</span></code> the row is a per-<code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">2</span></code> total (and <code class="docutils literal notranslate"><span class="pre">Type</span> <span class="pre">1</span></code> is <code class="docutils literal notranslate"><span class="pre">null</span></code>).</p> |
| </section> |
| </section> |
| <section id="aggregate-functions"> |
| <h2>Aggregate Functions<a class="headerlink" href="#aggregate-functions" title="Link to this heading">#</a></h2> |
| <p>The available aggregate functions are:</p> |
| <ol class="arabic simple"> |
| <li><dl class="simple myst"> |
| <dt>Comparison Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.min" title="datafusion.functions.min"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.min()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.max" title="datafusion.functions.max"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.max()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Math Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.sum" title="datafusion.functions.sum"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.sum()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.avg" title="datafusion.functions.avg"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.avg()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.median" title="datafusion.functions.median"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.median()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Array Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.array_agg" title="datafusion.functions.array_agg"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.array_agg()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Logical Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.bit_and" title="datafusion.functions.bit_and"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.bit_and()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.bit_or" title="datafusion.functions.bit_or"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.bit_or()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.bit_xor" title="datafusion.functions.bit_xor"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.bit_xor()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.bool_and" title="datafusion.functions.bool_and"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.bool_and()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.bool_or" title="datafusion.functions.bool_or"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.bool_or()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Statistical Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.count" title="datafusion.functions.count"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.count()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.corr" title="datafusion.functions.corr"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.corr()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.covar_samp" title="datafusion.functions.covar_samp"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.covar_samp()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.covar_pop" title="datafusion.functions.covar_pop"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.covar_pop()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.stddev" title="datafusion.functions.stddev"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.stddev()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.stddev_pop" title="datafusion.functions.stddev_pop"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.stddev_pop()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.var_samp" title="datafusion.functions.var_samp"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.var_samp()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.var_pop" title="datafusion.functions.var_pop"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.var_pop()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.var_population" title="datafusion.functions.var_population"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.var_population()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Linear Regression Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_count" title="datafusion.functions.regr_count"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_count()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_slope" title="datafusion.functions.regr_slope"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_slope()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_intercept" title="datafusion.functions.regr_intercept"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_intercept()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_r2" title="datafusion.functions.regr_r2"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_r2()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_avgx" title="datafusion.functions.regr_avgx"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_avgx()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_avgy" title="datafusion.functions.regr_avgy"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_avgy()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_sxx" title="datafusion.functions.regr_sxx"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_sxx()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_syy" title="datafusion.functions.regr_syy"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_syy()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.regr_slope" title="datafusion.functions.regr_slope"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.regr_slope()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Positional Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.first_value" title="datafusion.functions.first_value"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.first_value()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.last_value" title="datafusion.functions.last_value"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.last_value()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.nth_value" title="datafusion.functions.nth_value"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.nth_value()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>String Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.string_agg" title="datafusion.functions.string_agg"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.string_agg()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><dl class="simple myst"> |
| <dt>Percentile Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.percentile_cont" title="datafusion.functions.percentile_cont"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.percentile_cont()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.quantile_cont" title="datafusion.functions.quantile_cont"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.quantile_cont()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.approx_distinct" title="datafusion.functions.approx_distinct"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.approx_distinct()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.approx_median" title="datafusion.functions.approx_median"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.approx_median()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.approx_percentile_cont" title="datafusion.functions.approx_percentile_cont"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.approx_percentile_cont()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.approx_percentile_cont_with_weight" title="datafusion.functions.approx_percentile_cont_with_weight"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.approx_percentile_cont_with_weight()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| <li><p>Grouping Set Functions |
| - <a class="reference internal" href="../../autoapi/datafusion/functions/index.html#datafusion.functions.grouping" title="datafusion.functions.grouping"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.grouping()</span></code></a> |
| - <a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.rollup" title="datafusion.expr.GroupingSet.rollup"><code class="xref py py-meth docutils literal notranslate"><span class="pre">datafusion.expr.GroupingSet.rollup()</span></code></a> |
| - <a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.cube" title="datafusion.expr.GroupingSet.cube"><code class="xref py py-meth docutils literal notranslate"><span class="pre">datafusion.expr.GroupingSet.cube()</span></code></a> |
| - <a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.GroupingSet.grouping_sets" title="datafusion.expr.GroupingSet.grouping_sets"><code class="xref py py-meth docutils literal notranslate"><span class="pre">datafusion.expr.GroupingSet.grouping_sets()</span></code></a></p></li> |
| <li><dl class="simple myst"> |
| <dt>Spark-Compatible Functions</dt><dd><ul class="simple"> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/spark/index.html#datafusion.functions.spark.avg" title="datafusion.functions.spark.avg"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.spark.avg()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/spark/index.html#datafusion.functions.spark.try_sum" title="datafusion.functions.spark.try_sum"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.spark.try_sum()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/spark/index.html#datafusion.functions.spark.collect_list" title="datafusion.functions.spark.collect_list"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.spark.collect_list()</span></code></a></p></li> |
| <li><p><a class="reference internal" href="../../autoapi/datafusion/functions/spark/index.html#datafusion.functions.spark.collect_set" title="datafusion.functions.spark.collect_set"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.functions.spark.collect_set()</span></code></a></p></li> |
| </ul> |
| </dd> |
| </dl> |
| </li> |
| </ol> |
| <p>The functions in the <code class="docutils literal notranslate"><span class="pre">datafusion.functions.spark</span></code> namespace mirror Apache |
| Spark semantics, which can differ from the DataFusion built-ins of the same |
| name. They live in a separate namespace so you opt in explicitly. See |
| <a class="reference internal" href="spark-functions.html#spark-functions"><span class="std std-ref">Spark-Compatible Functions</span></a> for the full catalog and the semantic differences.</p> |
| </section> |
| <section id="user-defined-aggregate-functions"> |
| <h2>User-Defined Aggregate Functions<a class="headerlink" href="#user-defined-aggregate-functions" title="Link to this heading">#</a></h2> |
| <p>You can ship custom aggregations to the engine by subclassing |
| <a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#datafusion.user_defined.Accumulator" title="datafusion.user_defined.Accumulator"><code class="xref py py-class docutils literal notranslate"><span class="pre">Accumulator</span></code></a> and registering it via |
| <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.udaf" title="datafusion.udaf"><code class="xref py py-func docutils literal notranslate"><span class="pre">udaf()</span></code></a>. See <a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#module-datafusion.user_defined" title="datafusion.user_defined"><code class="xref py py-mod docutils literal notranslate"><span class="pre">datafusion.user_defined</span></code></a> for |
| the accumulator interface and worked examples.</p> |
| <div class="admonition note"> |
| <p class="admonition-title">Note</p> |
| <p>Serialization</p> |
| <p>Python aggregate UDFs travel inline inside pickled or |
| <a class="reference internal" href="../../autoapi/datafusion/expr/index.html#datafusion.expr.Expr.to_bytes" title="datafusion.expr.Expr.to_bytes"><code class="xref py py-meth docutils literal notranslate"><span class="pre">to_bytes()</span></code></a>-serialized expressions — |
| the accumulator class is captured by value via <code class="xref py py-mod docutils literal notranslate"><span class="pre">cloudpickle</span></code>, |
| so worker processes do not need to pre-register the UDF. Any names |
| the accumulator resolves via <code class="docutils literal notranslate"><span class="pre">import</span></code> are captured <strong>by reference</strong> |
| and must be importable on the receiving worker. See |
| <a class="reference internal" href="../../autoapi/datafusion/ipc/index.html#module-datafusion.ipc" title="datafusion.ipc"><code class="xref py py-mod docutils literal notranslate"><span class="pre">datafusion.ipc</span></code></a> for the full IPC model and security caveats.</p> |
| </div> |
| </section> |
| </section> |
| |
| |
| </article> |
| |
| |
| |
| |
| |
| <footer class="prev-next-footer d-print-none"> |
| |
| <div class="prev-next-area"> |
| <a class="left-prev" |
| href="spark-functions.html" |
| title="previous page"> |
| <i class="fa-solid fa-angle-left"></i> |
| <div class="prev-next-info"> |
| <p class="prev-next-subtitle">previous</p> |
| <p class="prev-next-title">Spark-Compatible Functions</p> |
| </div> |
| </a> |
| <a class="right-next" |
| href="windows.html" |
| title="next page"> |
| <div class="prev-next-info"> |
| <p class="prev-next-subtitle">next</p> |
| <p class="prev-next-title">Window Functions</p> |
| </div> |
| <i class="fa-solid fa-angle-right"></i> |
| </a> |
| </div> |
| </footer> |
| |
| </div> |
| |
| |
| |
| <dialog id="pst-secondary-sidebar-modal"></dialog> |
| <div id="pst-secondary-sidebar" class="bd-sidebar-secondary bd-toc"><div class="sidebar-secondary-items sidebar-secondary__inner"> |
| |
| |
| <div class="sidebar-secondary-item"> |
| <div |
| id="pst-page-navigation-heading-2" |
| class="page-toc tocsection onthispage"> |
| <i class="fa-solid fa-list"></i> On this page |
| </div> |
| <nav class="bd-toc-nav page-toc" aria-labelledby="pst-page-navigation-heading-2"> |
| <ul class="visible nav section-nav flex-column"> |
| <li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#setting-parameters">Setting Parameters</a><ul class="visible nav section-nav flex-column"> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#ordering">Ordering</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#distinct">Distinct</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#null-treatment">Null Treatment</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#filter">Filter</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#comparing-subsets-within-a-group">Comparing subsets within a group</a></li> |
| </ul> |
| </li> |
| <li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#grouping-sets">Grouping Sets</a><ul class="visible nav section-nav flex-column"> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#rollup">Rollup</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#cube">Cube</a></li> |
| <li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#explicit-grouping-sets">Explicit Grouping Sets</a></li> |
| </ul> |
| </li> |
| <li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#aggregate-functions">Aggregate Functions</a></li> |
| <li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#user-defined-aggregate-functions">User-Defined Aggregate Functions</a></li> |
| </ul> |
| </nav></div> |
| |
| </div></div> |
| |
| |
| </div> |
| <footer class="bd-footer-content"> |
| |
| </footer> |
| |
| </main> |
| </div> |
| </div> |
| |
| <!-- Scripts loaded after <body> so the DOM is not blocked --> |
| <script defer src="../../_static/scripts/bootstrap.js?digest=8878045cc6db502f8baf"></script> |
| <script defer src="../../_static/scripts/pydata-sphinx-theme.js?digest=8878045cc6db502f8baf"></script> |
| |
| <!-- Based on pydata_sphinx_theme/footer.html --> |
| <footer class="footer mt-5 mt-md-0"> |
| <div class="container"> |
| |
| <div class="footer-item"> |
| <p>Apache Arrow DataFusion, Arrow DataFusion, Apache, the Apache feather logo, and the Apache Arrow DataFusion project logo</p> |
| <p>are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.</p> |
| </div> |
| </div> |
| </footer> |
| |
| |
| </body> |
| </html> |