blob: 1db98b34366e1d68298969f398f90f61897e3911 [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>Execution Metrics &#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/dataframe/execution-metrics';</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="Common Operations" href="../common-operations/index.html" />
<link rel="prev" title="DataFrame Rendering" href="rendering.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 current active has-children"><a class="reference internal" href="index.html">DataFrames</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="rendering.html">DataFrame Rendering</a></li>
<li class="toctree-l3 current active"><a class="current reference internal" href="#">Execution Metrics</a></li>
</ul>
</details></li>
<li class="toctree-l2 has-children"><a class="reference internal" href="../common-operations/index.html">Common Operations</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="../common-operations/views.html">Registering Views</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/basic-info.html">Basic Operations</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/select-and-filter.html">Column Selections</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/expressions.html">Expressions</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/joins.html">Joins</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/functions.html">Functions</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/spark-functions.html">Spark-Compatible Functions</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/aggregations.html">Aggregation</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/windows.html">Window Functions</a></li>
<li class="toctree-l3"><a class="reference internal" href="../common-operations/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">DataFrames</a></li>
<li class="breadcrumb-item active" aria-current="page"><span class="ellipsis">Execution Metrics</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="execution-metrics">
<span id="id1"></span><h1>Execution Metrics<a class="headerlink" href="#execution-metrics" title="Link to this heading">#</a></h1>
<section id="overview">
<h2>Overview<a class="headerlink" href="#overview" title="Link to this heading">#</a></h2>
<p>When DataFusion executes a query it compiles the logical plan into a tree of
<em>physical plan operators</em> (e.g. <code class="docutils literal notranslate"><span class="pre">FilterExec</span></code>, <code class="docutils literal notranslate"><span class="pre">ProjectionExec</span></code>,
<code class="docutils literal notranslate"><span class="pre">HashAggregateExec</span></code>). Each operator can record runtime statistics while it
runs. These statistics are called <strong>execution metrics</strong>.</p>
<p>Typical metrics include:</p>
<ul class="simple">
<li><p><strong>output_rows</strong> – number of rows produced by the operator</p></li>
<li><p><strong>elapsed_compute</strong> – total CPU time (nanoseconds) spent inside the operator</p></li>
<li><p><strong>spill_count</strong> – number of times the operator spilled data to disk</p></li>
<li><p><strong>spilled_bytes</strong> – total bytes written to disk during spills</p></li>
<li><p><strong>spilled_rows</strong> – total rows written to disk during spills</p></li>
</ul>
<p>Metrics are collected <em>per-partition</em>: DataFusion may execute each operator
in parallel across several partitions. The convenience properties on
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">MetricsSet</span></code></a> (e.g. <code class="docutils literal notranslate"><span class="pre">output_rows</span></code>, <code class="docutils literal notranslate"><span class="pre">elapsed_compute</span></code>)
automatically sum the named metric across <strong>all</strong> partitions, giving a single
aggregate value for the operator as a whole. You can also access the raw
per-partition <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">Metric</span></code></a> objects via
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet.metrics" title="datafusion.MetricsSet.metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">metrics()</span></code></a>.</p>
</section>
<section id="when-are-metrics-available">
<h2>When Are Metrics Available?<a class="headerlink" href="#when-are-metrics-available" title="Link to this heading">#</a></h2>
<p>Some operators (for example <code class="docutils literal notranslate"><span class="pre">DataSourceExec</span></code>) eagerly create a
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">MetricsSet</span></code></a> when the physical plan is built, so
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan.metrics" title="datafusion.ExecutionPlan.metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">metrics()</span></code></a> may return a set even before any
rows have been processed. However, metric <strong>values</strong> such as <code class="docutils literal notranslate"><span class="pre">output_rows</span></code>
are only meaningful <strong>after</strong> the DataFrame has been executed via one of the
terminal operations:</p>
<ul class="simple">
<li><p><code class="xref py py-meth docutils literal notranslate"><span class="pre">collect()</span></code></p></li>
<li><p><code class="xref py py-meth docutils literal notranslate"><span class="pre">collect_partitioned()</span></code></p></li>
<li><p><code class="xref py py-meth docutils literal notranslate"><span class="pre">execute_stream()</span></code>
(metrics are available once the stream has been fully consumed)</p></li>
<li><p><code class="xref py py-meth docutils literal notranslate"><span class="pre">execute_stream_partitioned()</span></code>
(metrics are available once all partition streams have been fully consumed)</p></li>
</ul>
<p>Before execution, metric values will be <code class="docutils literal notranslate"><span class="pre">0</span></code> or <code class="docutils literal notranslate"><span class="pre">None</span></code>.</p>
<div class="admonition note">
<p class="admonition-title">Note</p>
<p><strong>display() does not populate metrics.</strong>
When a DataFrame is displayed in a notebook (e.g. via <code class="docutils literal notranslate"><span class="pre">display(df)</span></code> or
automatic <code class="docutils literal notranslate"><span class="pre">repr</span></code> output), DataFusion runs a <em>limited</em> internal execution
to fetch preview rows. This internal execution does <strong>not</strong> cache the
physical plan used, so <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan.collect_metrics" title="datafusion.ExecutionPlan.collect_metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">collect_metrics()</span></code></a>
will not reflect the display execution. To access metrics you must call
one of the terminal operations listed above.</p>
</div>
<p>If you call <code class="xref py py-meth docutils literal notranslate"><span class="pre">collect()</span></code> (or another terminal
operation) multiple times on the same DataFrame, each call creates a fresh
physical plan. Metrics from <code class="xref py py-meth docutils literal notranslate"><span class="pre">execution_plan()</span></code>
always reflect the <strong>most recent</strong> execution.</p>
</section>
<section id="reading-the-physical-plan-tree">
<h2>Reading the Physical Plan Tree<a class="headerlink" href="#reading-the-physical-plan-tree" title="Link to this heading">#</a></h2>
<p><code class="xref py py-meth docutils literal notranslate"><span class="pre">execution_plan()</span></code> returns the root
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan" title="datafusion.ExecutionPlan"><code class="xref py py-class docutils literal notranslate"><span class="pre">ExecutionPlan</span></code></a> node of the physical plan tree. The tree
mirrors the operator pipeline: the root is typically a projection or
coalescing node; its children are filters, aggregates, scans, etc.</p>
<p>The <code class="docutils literal notranslate"><span class="pre">operator_name</span></code> string returned by
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan.collect_metrics" title="datafusion.ExecutionPlan.collect_metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">collect_metrics()</span></code></a> is the <em>display</em> name of
the node, for example <code class="docutils literal notranslate"><span class="pre">&quot;FilterExec:</span> <span class="pre">column1&#64;0</span> <span class="pre">&gt;</span> <span class="pre">1&quot;</span></code>. This is the same string
you would see when calling <code class="docutils literal notranslate"><span class="pre">plan.display()</span></code>.</p>
</section>
<section id="aggregated-vs-per-partition-metrics">
<h2>Aggregated vs Per-Partition Metrics<a class="headerlink" href="#aggregated-vs-per-partition-metrics" title="Link to this heading">#</a></h2>
<p>DataFusion executes each operator across one or more <strong>partitions</strong> in
parallel. The <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">MetricsSet</span></code></a> convenience properties
(<code class="docutils literal notranslate"><span class="pre">output_rows</span></code>, <code class="docutils literal notranslate"><span class="pre">elapsed_compute</span></code>, etc.) automatically <strong>sum</strong> the named
metric across all partitions, giving a single aggregate value.</p>
<p>To inspect individual partitions — for example to detect data skew where one
partition processes far more rows than others — iterate over the raw
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">Metric</span></code></a> objects:</p>
<div class="highlight-python notranslate"><div class="highlight"><pre><span></span><span class="k">for</span> <span class="n">metric</span> <span class="ow">in</span> <span class="n">metrics_set</span><span class="o">.</span><span class="n">metrics</span><span class="p">():</span>
<span class="nb">print</span><span class="p">(</span><span class="sa">f</span><span class="s2">&quot; partition=</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">partition</span><span class="si">}</span><span class="s2"> </span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">name</span><span class="si">}</span><span class="s2">=</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">value</span><span class="si">}</span><span class="s2">&quot;</span><span class="p">)</span>
</pre></div>
</div>
<p>The <code class="docutils literal notranslate"><span class="pre">partition</span></code> property is a 0-based index (<code class="docutils literal notranslate"><span class="pre">0</span></code>, <code class="docutils literal notranslate"><span class="pre">1</span></code>, …) identifying
which parallel slot processed this metric. It is <code class="docutils literal notranslate"><span class="pre">None</span></code> for metrics that
apply globally (not tied to a specific partition).</p>
</section>
<section id="available-metrics">
<h2>Available Metrics<a class="headerlink" href="#available-metrics" title="Link to this heading">#</a></h2>
<p>The following metrics are directly accessible as properties on
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">MetricsSet</span></code></a>:</p>
<div class="pst-scrollable-table-container"><table class="table">
<thead>
<tr class="row-odd"><th class="head"><p>Property</p></th>
<th class="head"><p>Description</p></th>
</tr>
</thead>
<tbody>
<tr class="row-even"><td><p><code class="docutils literal notranslate"><span class="pre">output_rows</span></code></p></td>
<td><p>Number of rows emitted by the operator (summed across partitions).</p></td>
</tr>
<tr class="row-odd"><td><p><code class="docutils literal notranslate"><span class="pre">elapsed_compute</span></code></p></td>
<td><p>Wall-clock CPU time <strong>in nanoseconds</strong> spent inside the operator’s compute loop, excluding I/O wait. Useful for identifying which operators are most expensive (summed across partitions).</p></td>
</tr>
<tr class="row-even"><td><p><code class="docutils literal notranslate"><span class="pre">spill_count</span></code></p></td>
<td><p>Number of spill-to-disk events triggered by memory pressure. This is a unitless count of events, not a measure of data volume (summed across partitions).</p></td>
</tr>
<tr class="row-odd"><td><p><code class="docutils literal notranslate"><span class="pre">spilled_bytes</span></code></p></td>
<td><p>Total bytes written to disk during spill events (summed across partitions).</p></td>
</tr>
<tr class="row-even"><td><p><code class="docutils literal notranslate"><span class="pre">spilled_rows</span></code></p></td>
<td><p>Total rows written to disk during spill events (summed across partitions).</p></td>
</tr>
</tbody>
</table>
</div>
<p>Any metric not listed above can be accessed via
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet.sum_by_name" title="datafusion.MetricsSet.sum_by_name"><code class="xref py py-meth docutils literal notranslate"><span class="pre">sum_by_name()</span></code></a>, or by iterating over the raw
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">Metric</span></code></a> objects returned by
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet.metrics" title="datafusion.MetricsSet.metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">metrics()</span></code></a>.</p>
</section>
<section id="labels">
<h2>Labels<a class="headerlink" href="#labels" title="Link to this heading">#</a></h2>
<p>A <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">Metric</span></code></a> may carry <em>labels</em>: key/value pairs that
provide additional context. Labels are operator-specific; most metrics have
an empty label dict.</p>
<p>Some operators tag their metrics with labels to distinguish variants. For
example, a <code class="docutils literal notranslate"><span class="pre">HashAggregateExec</span></code> may record separate <code class="docutils literal notranslate"><span class="pre">output_rows</span></code> metrics
for intermediate and final output:</p>
<div class="highlight-python notranslate"><div class="highlight"><pre><span></span><span class="k">for</span> <span class="n">metric</span> <span class="ow">in</span> <span class="n">metrics_set</span><span class="o">.</span><span class="n">metrics</span><span class="p">():</span>
<span class="nb">print</span><span class="p">(</span><span class="n">metric</span><span class="o">.</span><span class="n">name</span><span class="p">,</span> <span class="n">metric</span><span class="o">.</span><span class="n">labels</span><span class="p">())</span>
<span class="c1"># output_rows {&#39;output_type&#39;: &#39;final&#39;}</span>
<span class="c1"># output_rows {&#39;output_type&#39;: &#39;intermediate&#39;}</span>
</pre></div>
</div>
<p>When summing by name (via <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet.output_rows" title="datafusion.MetricsSet.output_rows"><code class="xref py py-attr docutils literal notranslate"><span class="pre">output_rows</span></code></a> or
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet.sum_by_name" title="datafusion.MetricsSet.sum_by_name"><code class="xref py py-meth docutils literal notranslate"><span class="pre">sum_by_name()</span></code></a>), <strong>all</strong> metrics with that
name are summed regardless of labels. To filter by label, iterate over the
raw <a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">Metric</span></code></a> objects directly.</p>
</section>
<section id="end-to-end-example">
<h2>End-to-End Example<a class="headerlink" href="#end-to-end-example" title="Link to this heading">#</a></h2>
<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="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">sql</span><span class="p">(</span><span class="s2">&quot;CREATE TABLE sales AS VALUES (1, 100), (2, 200), (3, 50)&quot;</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">sql</span><span class="p">(</span><span class="s2">&quot;SELECT * FROM sales WHERE column1 &gt; 1&quot;</span><span class="p">)</span>
<span class="c1"># Execute the query — this populates the metrics</span>
<span class="n">results</span> <span class="o">=</span> <span class="n">df</span><span class="o">.</span><span class="n">collect</span><span class="p">()</span>
<span class="c1"># Retrieve the physical plan with metrics</span>
<span class="n">plan</span> <span class="o">=</span> <span class="n">df</span><span class="o">.</span><span class="n">execution_plan</span><span class="p">()</span>
<span class="c1"># Walk every operator and print its metrics</span>
<span class="k">for</span> <span class="n">operator_name</span><span class="p">,</span> <span class="n">ms</span> <span class="ow">in</span> <span class="n">plan</span><span class="o">.</span><span class="n">collect_metrics</span><span class="p">():</span>
<span class="k">if</span> <span class="n">ms</span><span class="o">.</span><span class="n">output_rows</span> <span class="ow">is</span> <span class="ow">not</span> <span class="kc">None</span><span class="p">:</span>
<span class="nb">print</span><span class="p">(</span><span class="sa">f</span><span class="s2">&quot;</span><span class="si">{</span><span class="n">operator_name</span><span class="si">}</span><span class="s2">&quot;</span><span class="p">)</span>
<span class="nb">print</span><span class="p">(</span><span class="sa">f</span><span class="s2">&quot; output_rows = </span><span class="si">{</span><span class="n">ms</span><span class="o">.</span><span class="n">output_rows</span><span class="si">}</span><span class="s2">&quot;</span><span class="p">)</span>
<span class="nb">print</span><span class="p">(</span><span class="sa">f</span><span class="s2">&quot; elapsed_compute = </span><span class="si">{</span><span class="n">ms</span><span class="o">.</span><span class="n">elapsed_compute</span><span class="si">}</span><span class="s2"> ns&quot;</span><span class="p">)</span>
<span class="c1"># Access raw per-partition metrics</span>
<span class="k">for</span> <span class="n">operator_name</span><span class="p">,</span> <span class="n">ms</span> <span class="ow">in</span> <span class="n">plan</span><span class="o">.</span><span class="n">collect_metrics</span><span class="p">():</span>
<span class="k">for</span> <span class="n">metric</span> <span class="ow">in</span> <span class="n">ms</span><span class="o">.</span><span class="n">metrics</span><span class="p">():</span>
<span class="nb">print</span><span class="p">(</span>
<span class="sa">f</span><span class="s2">&quot; partition=</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">partition</span><span class="si">}</span><span class="s2"> &quot;</span>
<span class="sa">f</span><span class="s2">&quot;</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">name</span><span class="si">}</span><span class="s2">=</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">value</span><span class="si">}</span><span class="s2"> &quot;</span>
<span class="sa">f</span><span class="s2">&quot;labels=</span><span class="si">{</span><span class="n">metric</span><span class="o">.</span><span class="n">labels</span><span class="p">()</span><span class="si">}</span><span class="s2">&quot;</span>
<span class="p">)</span>
</pre></div>
</div>
</section>
<section id="api-reference">
<h2>API Reference<a class="headerlink" href="#api-reference" title="Link to this heading">#</a></h2>
<ul class="simple">
<li><p><a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan" title="datafusion.ExecutionPlan"><code class="xref py py-class docutils literal notranslate"><span class="pre">datafusion.ExecutionPlan</span></code></a> — physical plan node</p></li>
<li><p><a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan.collect_metrics" title="datafusion.ExecutionPlan.collect_metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">datafusion.ExecutionPlan.collect_metrics()</span></code></a> — walk the tree and
return <code class="docutils literal notranslate"><span class="pre">(operator_name,</span> <span class="pre">MetricsSet)</span></code> pairs</p></li>
<li><p><a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.ExecutionPlan.metrics" title="datafusion.ExecutionPlan.metrics"><code class="xref py py-meth docutils literal notranslate"><span class="pre">datafusion.ExecutionPlan.metrics()</span></code></a> — return the
<a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">MetricsSet</span></code></a> for a single node</p></li>
<li><p><a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.MetricsSet" title="datafusion.MetricsSet"><code class="xref py py-class docutils literal notranslate"><span class="pre">datafusion.MetricsSet</span></code></a> — aggregated metrics for one operator</p></li>
<li><p><a class="reference internal" href="../../autoapi/datafusion/index.html#datafusion.Metric" title="datafusion.Metric"><code class="xref py py-class docutils literal notranslate"><span class="pre">datafusion.Metric</span></code></a> — a single per-partition metric value</p></li>
</ul>
</section>
</section>
</article>
<footer class="prev-next-footer d-print-none">
<div class="prev-next-area">
<a class="left-prev"
href="rendering.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">DataFrame Rendering</p>
</div>
</a>
<a class="right-next"
href="../common-operations/index.html"
title="next page">
<div class="prev-next-info">
<p class="prev-next-subtitle">next</p>
<p class="prev-next-title">Common Operations</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="#overview">Overview</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#when-are-metrics-available">When Are Metrics Available?</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#reading-the-physical-plan-tree">Reading the Physical Plan Tree</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#aggregated-vs-per-partition-metrics">Aggregated vs Per-Partition Metrics</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#available-metrics">Available Metrics</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#labels">Labels</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#end-to-end-example">End-to-End Example</a></li>
<li class="toc-h2 nav-item toc-entry"><a class="reference internal nav-link" href="#api-reference">API Reference</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>