blob: 196acf20e2a3d63da85fa8c5cbf8fe3e38f14f7a [file]
<!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>User-Defined Functions &#8212; 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/udf-and-udfa';</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="IO" href="../io/index.html" />
<link rel="prev" title="Window Functions" href="windows.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"><a class="reference internal" href="aggregations.html">Aggregation</a></li>
<li class="toctree-l3"><a class="reference internal" href="windows.html">Window Functions</a></li>
<li class="toctree-l3 current active"><a class="current reference internal" href="#">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">User-Defined Functions</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="user-defined-functions">
<h1>User-Defined Functions<a class="headerlink" href="#user-defined-functions" title="Link to this heading">#</a></h1>
<p>DataFusion provides powerful expressions and functions, reducing the need for custom Python
functions. However you can still incorporate your own functions, i.e. User-Defined Functions (UDFs).</p>
<section id="scalar-functions">
<h2>Scalar Functions<a class="headerlink" href="#scalar-functions" title="Link to this heading">#</a></h2>
<p>When writing a user-defined function that can operate on a row by row basis, these are called Scalar
Functions. You can define your own scalar function by calling
<a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#datafusion.user_defined.ScalarUDF.udf" title="datafusion.user_defined.ScalarUDF.udf"><code class="xref py py-func docutils literal notranslate"><span class="pre">udf()</span></code></a> .</p>
<p>The basic definition of a scalar UDF is a python function that takes one or more
<a class="reference external" href="https://arrow.apache.org/docs/python/index.html">pyarrow</a> arrays and returns a single array as
output. DataFusion scalar UDFs operate on an entire batch of records at a time, though the
evaluation of those records should be on a row by row basis. In the following example, we compute
if the input array contains null 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="kn">import</span><span class="w"> </span><span class="nn">pyarrow</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">datafusion</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">udf</span><span class="p">,</span> <span class="n">col</span>
<span class="k">def</span><span class="w"> </span><span class="nf">is_null</span><span class="p">(</span><span class="n">array</span><span class="p">:</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">Array</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">Array</span><span class="p">:</span>
<span class="k">return</span> <span class="n">array</span><span class="o">.</span><span class="n">is_null</span><span class="p">()</span>
<span class="n">is_null_arr</span> <span class="o">=</span> <span class="n">udf</span><span class="p">(</span><span class="n">is_null</span><span class="p">,</span> <span class="p">[</span><span class="n">pyarrow</span><span class="o">.</span><span class="n">int64</span><span class="p">()],</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">bool_</span><span class="p">(),</span> <span class="s1">&#39;stable&#39;</span><span class="p">)</span>
<span class="n">ctx</span> <span class="o">=</span> <span class="n">datafusion</span><span class="o">.</span><span class="n">SessionContext</span><span class="p">()</span>
<span class="n">batch</span> <span class="o">=</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">RecordBatch</span><span class="o">.</span><span class="n">from_arrays</span><span class="p">(</span>
<span class="p">[</span><span class="n">pyarrow</span><span class="o">.</span><span class="n">array</span><span class="p">([</span><span class="mi">1</span><span class="p">,</span> <span class="kc">None</span><span class="p">,</span> <span class="mi">3</span><span class="p">]),</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">array</span><span class="p">([</span><span class="mi">4</span><span class="p">,</span> <span class="mi">5</span><span class="p">,</span> <span class="mi">6</span><span class="p">])],</span>
<span class="n">names</span><span class="o">=</span><span class="p">[</span><span class="s2">&quot;a&quot;</span><span class="p">,</span> <span class="s2">&quot;b&quot;</span><span class="p">],</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">create_dataframe</span><span class="p">([[</span><span class="n">batch</span><span class="p">]],</span> <span class="n">name</span><span class="o">=</span><span class="s2">&quot;batch_array&quot;</span><span class="p">)</span>
<span class="n">df</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">&quot;a&quot;</span><span class="p">),</span> <span class="n">is_null_arr</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;a&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">&quot;is_null&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">show</span><span class="p">()</span>
</pre></div>
</div>
</div>
<div class="cell_output docutils container">
<div class="output stream highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame()
+---+---------+
| a | is_null |
+---+---------+
| 1 | false |
| | true |
| 3 | false |
+---+---------+
</pre></div>
</div>
</div>
</div>
<p>In the previous example, we used the fact that pyarrow provides a variety of built in array
functions such as <code class="docutils literal notranslate"><span class="pre">is_null()</span></code>. There are additional pyarrow
<a class="reference external" href="https://arrow.apache.org/docs/python/compute.html">compute functions</a> available. When possible,
it is highly recommended to use these functions because they can perform computations without doing
any copy operations from the original arrays. This leads to greatly improved performance.</p>
<p>If you need to perform an operation in python that is not available with the pyarrow compute
functions, you will need to convert the record batch into python values, perform your operation,
and construct an array. This operation of converting the built in data type of the array into a
python object can be one of the slowest operations in DataFusion, so it should be done sparingly.</p>
<p>The following example performs the same operation as before with <code class="docutils literal notranslate"><span class="pre">is_null</span></code> but demonstrates
converting to Python objects to do the evaluation.</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">import</span><span class="w"> </span><span class="nn">pyarrow</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">datafusion</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">udf</span><span class="p">,</span> <span class="n">col</span>
<span class="k">def</span><span class="w"> </span><span class="nf">is_null</span><span class="p">(</span><span class="n">array</span><span class="p">:</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">Array</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">Array</span><span class="p">:</span>
<span class="k">return</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">array</span><span class="p">([</span><span class="n">value</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span> <span class="ow">is</span> <span class="kc">None</span> <span class="k">for</span> <span class="n">value</span> <span class="ow">in</span> <span class="n">array</span><span class="p">])</span>
<span class="n">is_null_arr</span> <span class="o">=</span> <span class="n">udf</span><span class="p">(</span><span class="n">is_null</span><span class="p">,</span> <span class="p">[</span><span class="n">pyarrow</span><span class="o">.</span><span class="n">int64</span><span class="p">()],</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">bool_</span><span class="p">(),</span> <span class="s1">&#39;stable&#39;</span><span class="p">)</span>
<span class="n">ctx</span> <span class="o">=</span> <span class="n">datafusion</span><span class="o">.</span><span class="n">SessionContext</span><span class="p">()</span>
<span class="n">batch</span> <span class="o">=</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">RecordBatch</span><span class="o">.</span><span class="n">from_arrays</span><span class="p">(</span>
<span class="p">[</span><span class="n">pyarrow</span><span class="o">.</span><span class="n">array</span><span class="p">([</span><span class="mi">1</span><span class="p">,</span> <span class="kc">None</span><span class="p">,</span> <span class="mi">3</span><span class="p">]),</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">array</span><span class="p">([</span><span class="mi">4</span><span class="p">,</span> <span class="mi">5</span><span class="p">,</span> <span class="mi">6</span><span class="p">])],</span>
<span class="n">names</span><span class="o">=</span><span class="p">[</span><span class="s2">&quot;a&quot;</span><span class="p">,</span> <span class="s2">&quot;b&quot;</span><span class="p">],</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">create_dataframe</span><span class="p">([[</span><span class="n">batch</span><span class="p">]],</span> <span class="n">name</span><span class="o">=</span><span class="s2">&quot;batch_array&quot;</span><span class="p">)</span>
<span class="n">df</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">&quot;a&quot;</span><span class="p">),</span> <span class="n">is_null_arr</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;a&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">&quot;is_null&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">show</span><span class="p">()</span>
</pre></div>
</div>
</div>
<div class="cell_output docutils container">
<div class="output stream highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>DataFrame()
+---+---------+
| a | is_null |
+---+---------+
| 1 | false |
| | true |
| 3 | false |
+---+---------+
</pre></div>
</div>
</div>
</div>
<p>In this example we passed the PyArrow <code class="docutils literal notranslate"><span class="pre">DataType</span></code> when we defined the function
by calling <code class="docutils literal notranslate"><span class="pre">udf()</span></code>. If you need additional control, such as specifying
metadata or nullability of the input or output, you can instead specify a
PyArrow <code class="docutils literal notranslate"><span class="pre">Field</span></code>.</p>
<p>If you need to write a custom function but do not want to incur the performance
cost of converting to Python objects and back, a more advanced approach is to
write Rust based UDFs and to expose them to Python. There is an example in the
<a class="reference external" href="https://datafusion.apache.org/blog/2024/11/19/datafusion-python-udf-comparisons/">DataFusion blog</a>
describing how to do this.</p>
<section id="when-not-to-use-a-udf">
<h3>When not to use a UDF<a class="headerlink" href="#when-not-to-use-a-udf" title="Link to this heading">#</a></h3>
<p>A UDF is the right tool when the per-row computation genuinely cannot be
expressed with DataFusion’s built-in expressions. It is often the <em>wrong</em>
tool for a predicate that <em>can</em> be written as an <code class="docutils literal notranslate"><span class="pre">Expr</span></code> tree but feels
easier to write as a Python function — for example, a filter that keeps
a row if it matches any one of several rule sets, where each rule set
checks its own combination of columns (the worked example at the end of
this section keeps a row when it matches any one of several brand-specific
rules). Looping over the rules in Python and returning a boolean per row
reads naturally and is tempting to wrap in a UDF, but a UDF is opaque to
the optimizer: filters expressed as UDFs lose several rewrites that the
engine applies to filters built from native expressions. The most visible
of these is <strong>predicate pushdown into the table provider</strong>: a native
predicate can be handed to the source so it skips data before it is read,
while a UDF predicate cannot. The example below uses Parquet, where
pushdown prunes whole row groups using the min/max statistics in the
footer, but the same mechanism applies to any table provider that
advertises filter support — including custom providers.</p>
<p>The following example writes a small Parquet file, then filters it two
ways: first with a native expression, then with a UDF that computes the
same result. The filter itself is simple on purpose so we can compare
the plans side by side.</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">import</span><span class="w"> </span><span class="nn">tempfile</span><span class="o">,</span><span class="w"> </span><span class="nn">os</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">pa</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow.parquet</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">pq</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">udf</span>
<span class="n">tmpdir</span> <span class="o">=</span> <span class="n">tempfile</span><span class="o">.</span><span class="n">mkdtemp</span><span class="p">()</span>
<span class="n">parquet_path</span> <span class="o">=</span> <span class="n">os</span><span class="o">.</span><span class="n">path</span><span class="o">.</span><span class="n">join</span><span class="p">(</span><span class="n">tmpdir</span><span class="p">,</span> <span class="s2">&quot;items.parquet&quot;</span><span class="p">)</span>
<span class="n">pq</span><span class="o">.</span><span class="n">write_table</span><span class="p">(</span>
<span class="n">pa</span><span class="o">.</span><span class="n">table</span><span class="p">({</span>
<span class="s2">&quot;id&quot;</span><span class="p">:</span> <span class="nb">list</span><span class="p">(</span><span class="nb">range</span><span class="p">(</span><span class="mi">100</span><span class="p">)),</span>
<span class="s2">&quot;brand&quot;</span><span class="p">:</span> <span class="p">[</span><span class="s2">&quot;A&quot;</span><span class="p">,</span> <span class="s2">&quot;B&quot;</span><span class="p">,</span> <span class="s2">&quot;C&quot;</span><span class="p">,</span> <span class="s2">&quot;D&quot;</span><span class="p">]</span> <span class="o">*</span> <span class="mi">25</span><span class="p">,</span>
<span class="s2">&quot;qty&quot;</span><span class="p">:</span> <span class="p">[</span><span class="n">i</span> <span class="o">*</span> <span class="mi">10</span> <span class="k">for</span> <span class="n">i</span> <span class="ow">in</span> <span class="nb">range</span><span class="p">(</span><span class="mi">100</span><span class="p">)],</span>
<span class="p">}),</span>
<span class="n">parquet_path</span><span class="p">,</span>
<span class="p">)</span>
<span class="n">ctx</span> <span class="o">=</span> <span class="n">SessionContext</span><span class="p">()</span>
<span class="n">items</span> <span class="o">=</span> <span class="n">ctx</span><span class="o">.</span><span class="n">read_parquet</span><span class="p">(</span><span class="n">parquet_path</span><span class="p">)</span>
</pre></div>
</div>
</div>
</div>
<p><strong>Native-expression predicate.</strong> The filter is a plain boolean tree
over column references and literals, so the optimizer can analyze it:</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">native_filtered</span> <span class="o">=</span> <span class="n">items</span><span class="o">.</span><span class="n">filter</span><span class="p">(</span>
<span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;brand&quot;</span><span class="p">)</span> <span class="o">==</span> <span class="n">lit</span><span class="p">(</span><span class="s2">&quot;A&quot;</span><span class="p">))</span> <span class="o">&amp;</span> <span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;qty&quot;</span><span class="p">)</span> <span class="o">&gt;=</span> <span class="n">lit</span><span class="p">(</span><span class="mi">150</span><span class="p">))</span>
<span class="p">)</span>
<span class="nb">print</span><span class="p">(</span><span class="n">native_filtered</span><span class="o">.</span><span class="n">execution_plan</span><span class="p">()</span><span class="o">.</span><span class="n">display_indent</span><span class="p">())</span>
</pre></div>
</div>
</div>
<div class="cell_output docutils container">
<div class="output stream highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>FilterExec: brand@1 = A AND qty@2 &gt;= 150
RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
DataSourceExec: file_groups={1 group: [[tmp/tmp3q_xdwfo/items.parquet]]}, projection=[id, brand, qty], file_type=parquet, predicate=brand@1 = A AND qty@2 &gt;= 150, pruning_predicate=brand_null_count@2 != row_count@3 AND brand_min@0 &lt;= A AND A &lt;= brand_max@1 AND qty_null_count@5 != row_count@3 AND qty_max@4 &gt;= 150, required_guarantees=[brand in (A)]
</pre></div>
</div>
</div>
</div>
<p>Notice the <code class="docutils literal notranslate"><span class="pre">DataSourceExec</span></code> line. It carries three annotations the
optimizer computed from the predicate:</p>
<ul class="simple">
<li><p><code class="docutils literal notranslate"><span class="pre">predicate=brand&#64;1</span> <span class="pre">=</span> <span class="pre">A</span> <span class="pre">AND</span> <span class="pre">qty&#64;2</span> <span class="pre">&gt;=</span> <span class="pre">150</span></code> — the filter is pushed
into the Parquet scan itself, so the scan only reads matching rows.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">pruning_predicate=...</span> <span class="pre">brand_min&#64;0</span> <span class="pre">&lt;=</span> <span class="pre">A</span> <span class="pre">AND</span> <span class="pre">A</span> <span class="pre">&lt;=</span> <span class="pre">brand_max&#64;1</span> <span class="pre">...</span> <span class="pre">qty_max&#64;4</span> <span class="pre">&gt;=</span> <span class="pre">150</span></code> — the scan prunes whole row groups by consulting
the Parquet min/max statistics in the footer <em>before</em> reading any
column data.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">required_guarantees=[brand</span> <span class="pre">in</span> <span class="pre">(A)]</span></code> — the scan uses this when a
bloom filter or dictionary is available to skip pages.</p></li>
</ul>
<p><strong>UDF predicate.</strong> Now wrap the same logic in a Python UDF:</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="k">def</span><span class="w"> </span><span class="nf">brand_qty_filter</span><span class="p">(</span><span class="n">brand_arr</span><span class="p">:</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">,</span> <span class="n">qty_arr</span><span class="p">:</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">:</span>
<span class="k">return</span> <span class="n">pa</span><span class="o">.</span><span class="n">array</span><span class="p">([</span>
<span class="n">b</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span> <span class="o">==</span> <span class="s2">&quot;A&quot;</span> <span class="ow">and</span> <span class="n">q</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span> <span class="o">&gt;=</span> <span class="mi">150</span>
<span class="k">for</span> <span class="n">b</span><span class="p">,</span> <span class="n">q</span> <span class="ow">in</span> <span class="nb">zip</span><span class="p">(</span><span class="n">brand_arr</span><span class="p">,</span> <span class="n">qty_arr</span><span class="p">)</span>
<span class="p">])</span>
<span class="n">pred_udf</span> <span class="o">=</span> <span class="n">udf</span><span class="p">(</span>
<span class="n">brand_qty_filter</span><span class="p">,</span> <span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">string</span><span class="p">(),</span> <span class="n">pa</span><span class="o">.</span><span class="n">int64</span><span class="p">()],</span> <span class="n">pa</span><span class="o">.</span><span class="n">bool_</span><span class="p">(),</span> <span class="s2">&quot;stable&quot;</span><span class="p">,</span>
<span class="p">)</span>
<span class="n">udf_filtered</span> <span class="o">=</span> <span class="n">items</span><span class="o">.</span><span class="n">filter</span><span class="p">(</span><span class="n">pred_udf</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;brand&quot;</span><span class="p">),</span> <span class="n">col</span><span class="p">(</span><span class="s2">&quot;qty&quot;</span><span class="p">)))</span>
<span class="nb">print</span><span class="p">(</span><span class="n">udf_filtered</span><span class="o">.</span><span class="n">execution_plan</span><span class="p">()</span><span class="o">.</span><span class="n">display_indent</span><span class="p">())</span>
</pre></div>
</div>
</div>
<div class="cell_output docutils container">
<div class="output stream highlight-myst-ansi notranslate"><div class="highlight"><pre><span></span>FilterExec: brand_qty_filter(CAST(brand@1 AS Utf8), qty@2)
RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
DataSourceExec: file_groups={1 group: [[tmp/tmp3q_xdwfo/items.parquet]]}, projection=[id, brand, qty], file_type=parquet, predicate=brand_qty_filter(CAST(brand@1 AS Utf8), qty@2)
</pre></div>
</div>
</div>
</div>
<p>The <code class="docutils literal notranslate"><span class="pre">DataSourceExec</span></code> now carries only <code class="docutils literal notranslate"><span class="pre">predicate=brand_qty_filter(...)</span></code>.
There is no <code class="docutils literal notranslate"><span class="pre">pruning_predicate</span></code> and no <code class="docutils literal notranslate"><span class="pre">required_guarantees</span></code>: the
scan has to materialize every row group and hand each row to the
Python callback just to decide whether to keep it.</p>
<p>At small scale the cost difference is invisible; on a Parquet file with
many row groups, or data whose min/max statistics line up well with
the predicate, the native form can skip most of the file. The UDF form
reads all of it.</p>
<p><strong>Takeaway.</strong> Reach for a UDF when the per-row computation is genuinely
not expressible as a tree of built-in functions (custom numerical work,
external lookups, complex business rules). When it <em>is</em> expressible —
even if the native form is a little more verbose — build the <code class="docutils literal notranslate"><span class="pre">Expr</span></code>
tree directly so the optimizer can see through it. For disjunctive
predicates the idiom is to produce one clause per bucket and combine
them with <code class="docutils literal notranslate"><span class="pre">|</span></code>:</p>
<div class="highlight-python notranslate"><div class="highlight"><pre><span></span><span class="kn">from</span><span class="w"> </span><span class="nn">functools</span><span class="w"> </span><span class="kn">import</span> <span class="n">reduce</span>
<span class="kn">from</span><span class="w"> </span><span class="nn">operator</span><span class="w"> </span><span class="kn">import</span> <span class="n">or_</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">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">buckets</span> <span class="o">=</span> <span class="p">{</span>
<span class="s2">&quot;Brand#12&quot;</span><span class="p">:</span> <span class="p">{</span><span class="s2">&quot;containers&quot;</span><span class="p">:</span> <span class="p">[</span><span class="s2">&quot;SM CASE&quot;</span><span class="p">,</span> <span class="s2">&quot;SM BOX&quot;</span><span class="p">],</span> <span class="s2">&quot;min_qty&quot;</span><span class="p">:</span> <span class="mi">1</span><span class="p">,</span> <span class="s2">&quot;max_size&quot;</span><span class="p">:</span> <span class="mi">5</span><span class="p">},</span>
<span class="s2">&quot;Brand#23&quot;</span><span class="p">:</span> <span class="p">{</span><span class="s2">&quot;containers&quot;</span><span class="p">:</span> <span class="p">[</span><span class="s2">&quot;MED BAG&quot;</span><span class="p">,</span> <span class="s2">&quot;MED BOX&quot;</span><span class="p">],</span> <span class="s2">&quot;min_qty&quot;</span><span class="p">:</span> <span class="mi">10</span><span class="p">,</span> <span class="s2">&quot;max_size&quot;</span><span class="p">:</span> <span class="mi">10</span><span class="p">},</span>
<span class="p">}</span>
<span class="k">def</span><span class="w"> </span><span class="nf">bucket_clause</span><span class="p">(</span><span class="n">brand</span><span class="p">,</span> <span class="n">spec</span><span class="p">):</span>
<span class="k">return</span> <span class="p">(</span>
<span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;brand&quot;</span><span class="p">)</span> <span class="o">==</span> <span class="n">lit</span><span class="p">(</span><span class="n">brand</span><span class="p">))</span>
<span class="o">&amp;</span> <span class="n">f</span><span class="o">.</span><span class="n">in_list</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;container&quot;</span><span class="p">),</span> <span class="p">[</span><span class="n">lit</span><span class="p">(</span><span class="n">c</span><span class="p">)</span> <span class="k">for</span> <span class="n">c</span> <span class="ow">in</span> <span class="n">spec</span><span class="p">[</span><span class="s2">&quot;containers&quot;</span><span class="p">]])</span>
<span class="o">&amp;</span> <span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;quantity&quot;</span><span class="p">)</span> <span class="o">&gt;=</span> <span class="n">lit</span><span class="p">(</span><span class="n">spec</span><span class="p">[</span><span class="s2">&quot;min_qty&quot;</span><span class="p">]))</span>
<span class="o">&amp;</span> <span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;quantity&quot;</span><span class="p">)</span> <span class="o">&lt;=</span> <span class="n">lit</span><span class="p">(</span><span class="n">spec</span><span class="p">[</span><span class="s2">&quot;min_qty&quot;</span><span class="p">]</span> <span class="o">+</span> <span class="mi">10</span><span class="p">))</span>
<span class="o">&amp;</span> <span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;size&quot;</span><span class="p">)</span> <span class="o">&gt;=</span> <span class="n">lit</span><span class="p">(</span><span class="mi">1</span><span class="p">))</span>
<span class="o">&amp;</span> <span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;size&quot;</span><span class="p">)</span> <span class="o">&lt;=</span> <span class="n">lit</span><span class="p">(</span><span class="n">spec</span><span class="p">[</span><span class="s2">&quot;max_size&quot;</span><span class="p">]))</span>
<span class="p">)</span>
<span class="n">predicate</span> <span class="o">=</span> <span class="n">reduce</span><span class="p">(</span><span class="n">or_</span><span class="p">,</span> <span class="p">(</span><span class="n">bucket_clause</span><span class="p">(</span><span class="n">b</span><span class="p">,</span> <span class="n">s</span><span class="p">)</span> <span class="k">for</span> <span class="n">b</span><span class="p">,</span> <span class="n">s</span> <span class="ow">in</span> <span class="n">buckets</span><span class="o">.</span><span class="n">items</span><span class="p">()))</span>
<span class="n">df</span> <span class="o">=</span> <span class="n">df</span><span class="o">.</span><span class="n">filter</span><span class="p">(</span><span class="n">predicate</span><span class="p">)</span>
</pre></div>
</div>
</section>
</section>
<section id="aggregate-functions">
<h2>Aggregate Functions<a class="headerlink" href="#aggregate-functions" title="Link to this heading">#</a></h2>
<p>The <a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#datafusion.user_defined.AggregateUDF.udaf" title="datafusion.user_defined.AggregateUDF.udaf"><code class="xref py py-func docutils literal notranslate"><span class="pre">udaf()</span></code></a> function allows you to define User-Defined
Aggregate Functions (UDAFs). To use this you must implement an
<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> that determines how the aggregation is performed.</p>
<p>When defining a UDAF there are four methods you need to implement. The <code class="docutils literal notranslate"><span class="pre">update</span></code> function takes the
array(s) of input and updates the internal state of the accumulator. You should define this function
to have as many input arguments as you will pass when calling the UDAF. Since aggregation may be
split into multiple batches, we must have a method to combine multiple batches. For this, we have
two functions, <code class="docutils literal notranslate"><span class="pre">state</span></code> and <code class="docutils literal notranslate"><span class="pre">merge</span></code>. <code class="docutils literal notranslate"><span class="pre">state</span></code> will return an array of scalar values that contain
the current state of a single batch accumulation. Then we must <code class="docutils literal notranslate"><span class="pre">merge</span></code> the results of these
different states. Finally <code class="docutils literal notranslate"><span class="pre">evaluate</span></code> is the call that will return the final result after the
<code class="docutils literal notranslate"><span class="pre">merge</span></code> is complete.</p>
<p>In the following example we want to define a custom aggregate function that will return the
difference between the sum of two columns. The state can be represented by a single value and we can
also see how the inputs to <code class="docutils literal notranslate"><span class="pre">update</span></code> and <code class="docutils literal notranslate"><span class="pre">merge</span></code> differ.</p>
<div class="highlight-python notranslate"><div class="highlight"><pre><span></span><span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">pa</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow.compute</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">datafusion</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">col</span><span class="p">,</span> <span class="n">udaf</span><span class="p">,</span> <span class="n">Accumulator</span>
<span class="kn">from</span><span class="w"> </span><span class="nn">typing</span><span class="w"> </span><span class="kn">import</span> <span class="n">List</span>
<span class="k">class</span><span class="w"> </span><span class="nc">MyAccumulator</span><span class="p">(</span><span class="n">Accumulator</span><span class="p">):</span>
<span class="w"> </span><span class="sd">&quot;&quot;&quot;</span>
<span class="sd"> Interface of a user-defined accumulation.</span>
<span class="sd"> &quot;&quot;&quot;</span>
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_sum</span> <span class="o">=</span> <span class="mf">0.0</span>
<span class="k">def</span><span class="w"> </span><span class="nf">update</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">values_a</span><span class="p">:</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">,</span> <span class="n">values_b</span><span class="p">:</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_sum</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">_sum</span> <span class="o">+</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">compute</span><span class="o">.</span><span class="n">sum</span><span class="p">(</span><span class="n">values_a</span><span class="p">)</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span> <span class="o">-</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">compute</span><span class="o">.</span><span class="n">sum</span><span class="p">(</span><span class="n">values_b</span><span class="p">)</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span>
<span class="k">def</span><span class="w"> </span><span class="nf">merge</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">states</span><span class="p">:</span> <span class="nb">list</span><span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">])</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_sum</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">_sum</span> <span class="o">+</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">compute</span><span class="o">.</span><span class="n">sum</span><span class="p">(</span><span class="n">states</span><span class="p">[</span><span class="mi">0</span><span class="p">])</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span>
<span class="k">def</span><span class="w"> </span><span class="nf">state</span><span class="p">(</span><span class="bp">self</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="nb">list</span><span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">Scalar</span><span class="p">]:</span>
<span class="k">return</span> <span class="p">[</span><span class="n">pyarrow</span><span class="o">.</span><span class="n">scalar</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_sum</span><span class="p">)]</span>
<span class="k">def</span><span class="w"> </span><span class="nf">evaluate</span><span class="p">(</span><span class="bp">self</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">pa</span><span class="o">.</span><span class="n">Scalar</span><span class="p">:</span>
<span class="k">return</span> <span class="n">pyarrow</span><span class="o">.</span><span class="n">scalar</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_sum</span><span class="p">)</span>
<span class="n">ctx</span> <span class="o">=</span> <span class="n">datafusion</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">from_pydict</span><span class="p">(</span>
<span class="p">{</span>
<span class="s2">&quot;a&quot;</span><span class="p">:</span> <span class="p">[</span><span class="mi">4</span><span class="p">,</span> <span class="mi">5</span><span class="p">,</span> <span class="mi">6</span><span class="p">],</span>
<span class="s2">&quot;b&quot;</span><span class="p">:</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">3</span><span class="p">],</span>
<span class="p">}</span>
<span class="p">)</span>
<span class="n">my_udaf</span> <span class="o">=</span> <span class="n">udaf</span><span class="p">(</span><span class="n">MyAccumulator</span><span class="p">,</span> <span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">(),</span> <span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">()],</span> <span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">(),</span> <span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">()],</span> <span class="s1">&#39;stable&#39;</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="p">[</span><span class="n">my_udaf</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;a&quot;</span><span class="p">),</span> <span class="n">col</span><span class="p">(</span><span class="s2">&quot;b&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">&quot;col_diff&quot;</span><span class="p">)])</span>
</pre></div>
</div>
<section id="faq">
<h3>FAQ<a class="headerlink" href="#faq" title="Link to this heading">#</a></h3>
<p><strong>How do I return a list from a UDAF?</strong></p>
<p>Both the <code class="docutils literal notranslate"><span class="pre">evaluate</span></code> and the <code class="docutils literal notranslate"><span class="pre">state</span></code> functions expect to return scalar values.
If you wish to return a list array as a scalar value, the best practice is to
wrap the values in a <code class="docutils literal notranslate"><span class="pre">pyarrow.Scalar</span></code> object. For example, you can return a
timestamp list with <code class="docutils literal notranslate"><span class="pre">pa.scalar([...],</span> <span class="pre">type=pa.list_(pa.timestamp(&quot;ms&quot;)))</span></code> and
register the appropriate return or state types as
<code class="docutils literal notranslate"><span class="pre">return_type=pa.list_(pa.timestamp(&quot;ms&quot;))</span></code> and
<code class="docutils literal notranslate"><span class="pre">state_type=[pa.list_(pa.timestamp(&quot;ms&quot;))]</span></code>, respectively.</p>
<p>As of DataFusion 52.0.0 , you can pass return any Python object, including a
PyArrow array, as the return value(s) for these functions and DataFusion will
attempt to create a scalar type from the value. DataFusion has been tested to
convert PyArrow, nanoarrow, and arro3 objects as well as primitive data types
like integers, strings, and so on.</p>
</section>
</section>
<section id="window-functions">
<h2>Window Functions<a class="headerlink" href="#window-functions" title="Link to this heading">#</a></h2>
<p>To implement a User-Defined Window Function (UDWF) you must call the
<a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#datafusion.user_defined.WindowUDF.udwf" title="datafusion.user_defined.WindowUDF.udwf"><code class="xref py py-func docutils literal notranslate"><span class="pre">udwf()</span></code></a> function using a class that implements the abstract
class <a class="reference internal" href="../../autoapi/datafusion/user_defined/index.html#datafusion.user_defined.WindowEvaluator" title="datafusion.user_defined.WindowEvaluator"><code class="xref py py-class docutils literal notranslate"><span class="pre">WindowEvaluator</span></code></a>.</p>
<p>There are three methods of evaluation of UDWFs.</p>
<ul class="simple">
<li><p><code class="docutils literal notranslate"><span class="pre">evaluate</span></code> is the simplest case, where you are given an array and are expected to calculate the
value for a single row of that array. This is the simplest case, but also the least performant.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">evaluate_all</span></code> computes the values for all rows for an input array at a single time.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">evaluate_all_with_rank</span></code> computes the values for all rows, but you only have the rank
information for the rows.</p></li>
</ul>
<p>Which methods you implement are based upon which of these options are set.</p>
<div class="pst-scrollable-table-container"><table class="table">
<thead>
<tr class="row-odd"><th class="head"><p><code class="docutils literal notranslate"><span class="pre">uses_window_frame</span></code></p></th>
<th class="head"><p><code class="docutils literal notranslate"><span class="pre">supports_bounded_execution</span></code></p></th>
<th class="head"><p><code class="docutils literal notranslate"><span class="pre">include_rank</span></code></p></th>
<th class="head"><p>function_to_implement</p></th>
</tr>
</thead>
<tbody>
<tr class="row-even"><td><p>False (default)</p></td>
<td><p>False (default)</p></td>
<td><p>False (default)</p></td>
<td><p><code class="docutils literal notranslate"><span class="pre">evaluate_all</span></code></p></td>
</tr>
<tr class="row-odd"><td><p>False</p></td>
<td><p>True</p></td>
<td><p>False</p></td>
<td><p><code class="docutils literal notranslate"><span class="pre">evaluate</span></code></p></td>
</tr>
<tr class="row-even"><td><p>False</p></td>
<td><p>True</p></td>
<td><p>False</p></td>
<td><p><code class="docutils literal notranslate"><span class="pre">evaluate_all_with_rank</span></code></p></td>
</tr>
<tr class="row-odd"><td><p>True</p></td>
<td><p>True/False</p></td>
<td><p>True/False</p></td>
<td><p><code class="docutils literal notranslate"><span class="pre">evaluate</span></code></p></td>
</tr>
</tbody>
</table>
</div>
<section id="udwf-options">
<h3>UDWF options<a class="headerlink" href="#udwf-options" title="Link to this heading">#</a></h3>
<p>When you define your UDWF you can override the functions that return these values. They will
determine which evaluate functions are called.</p>
<ul class="simple">
<li><p><code class="docutils literal notranslate"><span class="pre">uses_window_frame</span></code> is set for functions that compute based on the specified window frame. If
your function depends upon the specified frame, set this to <code class="docutils literal notranslate"><span class="pre">True</span></code>.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">supports_bounded_execution</span></code> specifies if your function can be incrementally computed.</p></li>
<li><p><code class="docutils literal notranslate"><span class="pre">include_rank</span></code> is set to <code class="docutils literal notranslate"><span class="pre">True</span></code> for window functions that can be computed only using the rank
information.</p></li>
</ul>
<div class="highlight-python notranslate"><div class="highlight"><pre><span></span><span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">pa</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">udwf</span><span class="p">,</span> <span class="n">col</span><span class="p">,</span> <span class="n">SessionContext</span>
<span class="kn">from</span><span class="w"> </span><span class="nn">datafusion.user_defined</span><span class="w"> </span><span class="kn">import</span> <span class="n">WindowEvaluator</span>
<span class="k">class</span><span class="w"> </span><span class="nc">ExponentialSmooth</span><span class="p">(</span><span class="n">WindowEvaluator</span><span class="p">):</span>
<span class="k">def</span><span class="w"> </span><span class="fm">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">alpha</span><span class="p">:</span> <span class="nb">float</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="kc">None</span><span class="p">:</span>
<span class="bp">self</span><span class="o">.</span><span class="n">alpha</span> <span class="o">=</span> <span class="n">alpha</span>
<span class="k">def</span><span class="w"> </span><span class="nf">evaluate_all</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">values</span><span class="p">:</span> <span class="nb">list</span><span class="p">[</span><span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">],</span> <span class="n">num_rows</span><span class="p">:</span> <span class="nb">int</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">pa</span><span class="o">.</span><span class="n">Array</span><span class="p">:</span>
<span class="n">results</span> <span class="o">=</span> <span class="p">[]</span>
<span class="n">curr_value</span> <span class="o">=</span> <span class="mf">0.0</span>
<span class="n">values</span> <span class="o">=</span> <span class="n">values</span><span class="p">[</span><span class="mi">0</span><span class="p">]</span>
<span class="k">for</span> <span class="n">idx</span> <span class="ow">in</span> <span class="nb">range</span><span class="p">(</span><span class="n">num_rows</span><span class="p">):</span>
<span class="k">if</span> <span class="n">idx</span> <span class="o">==</span> <span class="mi">0</span><span class="p">:</span>
<span class="n">curr_value</span> <span class="o">=</span> <span class="n">values</span><span class="p">[</span><span class="n">idx</span><span class="p">]</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span>
<span class="k">else</span><span class="p">:</span>
<span class="n">curr_value</span> <span class="o">=</span> <span class="n">values</span><span class="p">[</span><span class="n">idx</span><span class="p">]</span><span class="o">.</span><span class="n">as_py</span><span class="p">()</span> <span class="o">*</span> <span class="bp">self</span><span class="o">.</span><span class="n">alpha</span> <span class="o">+</span> <span class="n">curr_value</span> <span class="o">*</span> <span class="p">(</span>
<span class="mf">1.0</span> <span class="o">-</span> <span class="bp">self</span><span class="o">.</span><span class="n">alpha</span>
<span class="p">)</span>
<span class="n">results</span><span class="o">.</span><span class="n">append</span><span class="p">(</span><span class="n">curr_value</span><span class="p">)</span>
<span class="k">return</span> <span class="n">pa</span><span class="o">.</span><span class="n">array</span><span class="p">(</span><span class="n">results</span><span class="p">)</span>
<span class="n">exp_smooth</span> <span class="o">=</span> <span class="n">udwf</span><span class="p">(</span>
<span class="n">ExponentialSmooth</span><span class="p">(</span><span class="mf">0.9</span><span class="p">),</span>
<span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">(),</span>
<span class="n">pa</span><span class="o">.</span><span class="n">float64</span><span class="p">(),</span>
<span class="n">volatility</span><span class="o">=</span><span class="s2">&quot;immutable&quot;</span><span class="p">,</span>
<span class="p">)</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">from_pydict</span><span class="p">({</span>
<span class="s2">&quot;a&quot;</span><span class="p">:</span> <span class="p">[</span><span class="mf">1.0</span><span class="p">,</span> <span class="mf">2.1</span><span class="p">,</span> <span class="mf">2.9</span><span class="p">,</span> <span class="mf">4.0</span><span class="p">,</span> <span class="mf">5.1</span><span class="p">,</span> <span class="mf">6.0</span><span class="p">,</span> <span class="mf">6.9</span><span class="p">,</span> <span class="mf">8.0</span><span class="p">]</span>
<span class="p">})</span>
<span class="n">df</span><span class="o">.</span><span class="n">select</span><span class="p">(</span><span class="s2">&quot;a&quot;</span><span class="p">,</span> <span class="n">exp_smooth</span><span class="p">(</span><span class="n">col</span><span class="p">(</span><span class="s2">&quot;a&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">alias</span><span class="p">(</span><span class="s2">&quot;smooth_a&quot;</span><span class="p">))</span><span class="o">.</span><span class="n">show</span><span class="p">()</span>
</pre></div>
</div>
</section>
</section>
<section id="table-functions">
<h2>Table Functions<a class="headerlink" href="#table-functions" title="Link to this heading">#</a></h2>
<p>User Defined Table Functions are slightly different than the other functions
described here. These functions take any number of <code class="docutils literal notranslate"><span class="pre">Expr</span></code> arguments, but only
literal expressions are supported. Table functions must return a Table
Provider as described in the ref:<code class="docutils literal notranslate"><span class="pre">_io_custom_table_provider</span></code> page.</p>
<p>Once you have a table function, you can register it with the session context
by using <a class="reference internal" href="../../autoapi/datafusion/context/index.html#datafusion.context.SessionContext.register_udtf" title="datafusion.context.SessionContext.register_udtf"><code class="xref py py-func docutils literal notranslate"><span class="pre">datafusion.context.SessionContext.register_udtf()</span></code></a>.</p>
<p>There are examples of both rust backed and python based table functions in the
examples folder of the repository. If you have a rust backed table function
that you wish to expose via PyO3, you need to expose it as a <code class="docutils literal notranslate"><span class="pre">PyCapsule</span></code>.</p>
<div class="highlight-rust notranslate"><div class="highlight"><pre><span></span><span class="cp">#[pymethods]</span>
<span class="k">impl</span><span class="w"> </span><span class="n">MyTableFunction</span><span class="w"> </span><span class="p">{</span>
<span class="w"> </span><span class="k">fn</span><span class="w"> </span><span class="nf">__datafusion_table_function__</span><span class="o">&lt;&#39;</span><span class="na">py</span><span class="o">&gt;</span><span class="p">(</span>
<span class="w"> </span><span class="o">&amp;</span><span class="bp">self</span><span class="p">,</span>
<span class="w"> </span><span class="n">py</span><span class="p">:</span><span class="w"> </span><span class="nc">Python</span><span class="o">&lt;&#39;</span><span class="na">py</span><span class="o">&gt;</span><span class="p">,</span>
<span class="w"> </span><span class="p">)</span><span class="w"> </span><span class="p">-&gt;</span><span class="w"> </span><span class="nc">PyResult</span><span class="o">&lt;</span><span class="n">Bound</span><span class="o">&lt;&#39;</span><span class="na">py</span><span class="p">,</span><span class="w"> </span><span class="n">PyCapsule</span><span class="o">&gt;&gt;</span><span class="w"> </span><span class="p">{</span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">name</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="n">cr</span><span class="s">&quot;datafusion_table_function&quot;</span><span class="p">.</span><span class="n">into</span><span class="p">();</span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">func</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="bp">self</span><span class="p">.</span><span class="n">clone</span><span class="p">();</span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">provider</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="n">FFI_TableFunction</span><span class="p">::</span><span class="n">new</span><span class="p">(</span><span class="n">Arc</span><span class="p">::</span><span class="n">new</span><span class="p">(</span><span class="n">func</span><span class="p">),</span><span class="w"> </span><span class="nb">None</span><span class="p">);</span>
<span class="w"> </span><span class="n">PyCapsule</span><span class="p">::</span><span class="n">new</span><span class="p">(</span><span class="n">py</span><span class="p">,</span><span class="w"> </span><span class="n">provider</span><span class="p">,</span><span class="w"> </span><span class="nb">Some</span><span class="p">(</span><span class="n">name</span><span class="p">))</span>
<span class="w"> </span><span class="p">}</span>
<span class="p">}</span>
</pre></div>
</div>
<section id="accessing-the-calling-session">
<h3>Accessing the Calling Session<a class="headerlink" href="#accessing-the-calling-session" title="Link to this heading">#</a></h3>
<p>Pure-Python UDTFs can opt into receiving the calling
<code class="xref py py-class docutils literal notranslate"><span class="pre">SessionContext</span></code> by registering with
<code class="docutils literal notranslate"><span class="pre">with_session=True</span></code>. The context is passed as a <code class="docutils literal notranslate"><span class="pre">session</span></code> keyword
argument on every invocation. Use it to look up registered tables,
UDFs, or session configuration from inside the callback.</p>
<div class="highlight-python 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">Table</span><span class="p">,</span> <span class="n">udtf</span>
<span class="kn">from</span><span class="w"> </span><span class="nn">datafusion.context</span><span class="w"> </span><span class="kn">import</span> <span class="n">TableProviderExportable</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">pa</span>
<span class="kn">import</span><span class="w"> </span><span class="nn">pyarrow.dataset</span><span class="w"> </span><span class="k">as</span><span class="w"> </span><span class="nn">ds</span>
<span class="nd">@udtf</span><span class="p">(</span><span class="s2">&quot;list_tables&quot;</span><span class="p">,</span> <span class="n">with_session</span><span class="o">=</span><span class="kc">True</span><span class="p">)</span>
<span class="k">def</span><span class="w"> </span><span class="nf">list_tables</span><span class="p">(</span><span class="o">*</span><span class="p">,</span> <span class="n">session</span><span class="p">:</span> <span class="n">SessionContext</span><span class="p">)</span> <span class="o">-&gt;</span> <span class="n">TableProviderExportable</span><span class="p">:</span>
<span class="n">names</span> <span class="o">=</span> <span class="nb">sorted</span><span class="p">(</span><span class="n">session</span><span class="o">.</span><span class="n">catalog</span><span class="p">()</span><span class="o">.</span><span class="n">schema</span><span class="p">()</span><span class="o">.</span><span class="n">names</span><span class="p">())</span>
<span class="n">batch</span> <span class="o">=</span> <span class="n">pa</span><span class="o">.</span><span class="n">RecordBatch</span><span class="o">.</span><span class="n">from_pydict</span><span class="p">({</span><span class="s2">&quot;name&quot;</span><span class="p">:</span> <span class="n">names</span><span class="p">})</span>
<span class="k">return</span> <span class="n">Table</span><span class="p">(</span><span class="n">ds</span><span class="o">.</span><span class="n">dataset</span><span class="p">([</span><span class="n">batch</span><span class="p">]))</span>
<span class="n">ctx</span> <span class="o">=</span> <span class="n">SessionContext</span><span class="p">()</span>
<span class="n">ctx</span><span class="o">.</span><span class="n">register_batch</span><span class="p">(</span><span class="s2">&quot;t1&quot;</span><span class="p">,</span> <span class="n">pa</span><span class="o">.</span><span class="n">RecordBatch</span><span class="o">.</span><span class="n">from_pydict</span><span class="p">({</span><span class="s2">&quot;x&quot;</span><span class="p">:</span> <span class="p">[</span><span class="mi">1</span><span class="p">]}))</span>
<span class="n">ctx</span><span class="o">.</span><span class="n">register_udtf</span><span class="p">(</span><span class="n">list_tables</span><span class="p">)</span>
<span class="n">ctx</span><span class="o">.</span><span class="n">sql</span><span class="p">(</span><span class="s2">&quot;SELECT * FROM list_tables()&quot;</span><span class="p">)</span><span class="o">.</span><span class="n">show</span><span class="p">()</span>
</pre></div>
</div>
<p>Without <code class="docutils literal notranslate"><span class="pre">with_session=True</span></code>, the callback receives only the positional
expression arguments. The flag is opt-in so existing UDTFs keep working
unchanged.</p>
<p>The injected <code class="docutils literal notranslate"><span class="pre">session</span></code> is a fresh <code class="xref py py-class docutils literal notranslate"><span class="pre">SessionContext</span></code>
wrapper backed by the same underlying state as the caller, so registries
(tables, UDFs, catalogs) are visible. Registry mutations (e.g. registering
a new table or UDF) propagate to the live session because the registries
are reference-counted and shared. Configuration changes made through the
wrapper (e.g. setting session options) do <strong>not</strong> propagate — the wrapper
holds its own clone of the session config.</p>
</section>
</section>
</section>
</article>
<footer class="prev-next-footer d-print-none">
<div class="prev-next-area">
<a class="left-prev"
href="windows.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">Window Functions</p>
</div>
</a>
<a class="right-next"
href="../io/index.html"
title="next page">
<div class="prev-next-info">
<p class="prev-next-subtitle">next</p>
<p class="prev-next-title">IO</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="#scalar-functions">Scalar Functions</a><ul class="visible nav section-nav flex-column">
<li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#when-not-to-use-a-udf">When not to use a UDF</a></li>
</ul>
</li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#aggregate-functions">Aggregate Functions</a><ul class="visible nav section-nav flex-column">
<li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#faq">FAQ</a></li>
</ul>
</li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#window-functions">Window Functions</a><ul class="visible nav section-nav flex-column">
<li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#udwf-options">UDWF options</a></li>
</ul>
</li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#table-functions">Table Functions</a><ul class="visible nav section-nav flex-column">
<li class="toc-h3 nav-item toc-entry"><a class="reference internal nav-link" href="#accessing-the-calling-session">Accessing the Calling Session</a></li>
</ul>
</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>