Skip to content

DataFrame Operations

Core row- and column-level transforms: filter, select, add or modify columns, sort, deduplicate, and the string/date/list expression namespaces. Every FlowFrame method mirrors its Polars counterpart and accepts an optional description that shows up as node documentation in the visual editor.

The flowfile_formula examples below use the Flowfile formula language; everything else is a Polars expression.

A worked example

This runs against committed data and is executed by the docs test suite:

import flowfile as ff

# Deduplicate, derive columns, filter, sort, then drop a column with a selector.
cleaned = (
    ff.read_csv("data/templates/orders.csv")
    .unique(subset=["order_id"])
    .with_columns(
        (ff.col("quantity") * ff.lit(2)).alias("double_qty"),
        ff.col("quantity").cast(ff.Float64).alias("quantity_f"),
    )
    .filter(ff.col("quantity") >= 3)
    .sort("quantity", descending=True)
    .select(ff.col("*").exclude("product_id"))
    .collect()
)

# String and date namespaces on expressions.
enriched = (
    ff.read_csv("data/templates/customers.csv")
    .with_columns(
        ff.col("name").str.to_uppercase().alias("name_upper"),
        ff.col("email").str.slice(0, 5).alias("email_prefix"),
        ff.col("name").str.contains("a").alias("name_has_a"),
        ff.col("signup_date").dt.year().alias("signup_year"),
        ff.col("signup_date").dt.weekday().alias("signup_weekday"),
    )
    .collect()
)

# A running total per customer with a window (cum_sum over a group).
running = (
    ff.read_csv("data/templates/orders.csv")
    .sort("order_date")
    .with_columns(ff.col("quantity").cum_sum().over("customer_id").alias("running_qty"))
    .collect()
)

The sections below break down each operation.

Filtering

import flowfile as ff

df = ff.FlowFrame({"price": [10, 20, 30], "qty": [5, 0, 10]})

# Polars expression predicate
df = df.filter(ff.col("price") > 15)

# With a description (surfaces in the visual editor)
df = df.filter(ff.col("price") > 15, description="Keep items over $15")

# Flowfile formula syntax
df = df.filter(flowfile_formula="[price] > 15 and [qty] > 0")

Which node the filter becomes

Both filter(flowfile_formula=...) and a plain filter(ff.col(...) > x) predicate emit an editable Filter node. A predicate with no formula form, such as a lambda, falls back to a polars_code node; see which operations become which node.

Selecting columns

# Select specific columns by name
df = df.select(["price", "qty"])

# Select with expressions
df = df.select([
    ff.col("price"),
    ff.col("qty").alias("quantity"),
])

# Keep everything except one column with a column selector.
# There is no ff.exclude() — use ff.col("*").exclude(...) or the selectors module.
df = df.select(ff.col("*").exclude("internal_id"))

Selectors

ff.numeric(), ff.string(), ff.all_(), and the other selector helpers pick columns by dtype or pattern — e.g. df.select(ff.numeric()) keeps only numeric columns.

Adding and modifying columns

# Expression form
df = df.with_columns([
    (ff.col("price") * ff.col("qty")).alias("total"),
])

# Flowfile formula form
df = df.with_columns(
    flowfile_formulas=["[price] * [qty]"],
    output_column_names=["total"],
    description="Calculate line totals",
)

Sorting

df = df.sort("price")
df = df.sort("price", descending=True)

# Multi-column sort
df = df.sort(["category", "price"], descending=[False, True])

Removing duplicates

# Drop fully duplicate rows
df = df.unique()

# Deduplicate on a subset of columns
df = df.unique(subset=["product_id"])

No drop_duplicates

FlowFrame does not expose drop_duplicates. Use unique() (optionally with subset=[...] and keep="first").

Cleaning messy text

data_cleansing() is the Python form of the Data Cleansing node: one call that fills nulls, trims and normalizes whitespace, strips unwanted characters, and fixes casing. Pass a list of column names to limit it; with no list it cleanses every column. Text rules only touch String columns and the null-to-zero rule only Numeric ones, so a mixed frame is safe to pass whole.

import flowfile as ff

raw = ff.from_dict(
    {
        "name": ["  alice smith", "BOB   JONES ", None, None],
        "city": ["amsterdam", " Rotterdam  ", "utrecht", None],
        "score": [10, None, None, None],
    }
)

# Defaults already fill nulls (blank for text, 0 for numbers) and trim whitespace.
cleaned = raw.data_cleansing(
    remove_null_rows=True,
    normalize_whitespace=True,
    case_mode="titlecase",
).collect()

# Character rules apply only to the columns you name; other columns pass through.
phones = (
    ff.from_dict({"phone": ["+31 (0)20-123 4567", "06 1234 5678"], "id": [1, 2]})
    .data_cleansing(["phone"], remove_punctuation=True, remove_all_whitespace=True)
    .collect()
)

remove_null_rows and remove_null_columns look at the whole frame regardless of the column list. remove_null_columns is the one data-dependent keyword: which columns go is decided from a null count taken when the method is called, not at collect().

One formula over many columns

multi_field_formula() is the Python form of the Multi-Field Formula node: one Flowfile formula applied to many columns in a single call. Inside the formula [_CurrentField_] is the value of the column being processed, [_CurrentFieldName_] its name and [_CurrentFieldType_] its dtype, the last two as text. Narrow the selection with columns=[...] or data_type="Numeric" — they are mutually exclusive, and passing neither applies the formula to every column.

import flowfile as ff

budget = ff.from_dict(
    {
        "region": ["north", "south"],
        "jan": [1200, 900],
        "feb": [1500, 1100],
        "total": [2700, 2000],
    }
)

# [_CurrentField_] is the value of the column being processed. A prefix or a suffix
# writes the results to new columns and leaves the source columns untouched.
with_vat = budget.multi_field_formula(
    "round([_CurrentField_] * 1.21, 2)",
    data_type="Numeric",
    suffix="_incl_vat",
    output_data_type="Float64",
).collect()

# Without an affix the results overwrite their source columns. Every expression reads
# the original input, so [total] is still the input total while `total` is overwritten.
percent = budget.multi_field_formula(
    "round([_CurrentField_] / [total] * 100, 1)",
    columns=["total", "jan", "feb"],
).collect()

# [_CurrentFieldName_] and [_CurrentFieldType_] bind the column's name and its dtype.
labelled = budget.multi_field_formula(
    'concat([_CurrentFieldName_], " is ", [_CurrentFieldType_])',
    data_type="String",
    prefix="about_",
).collect()

prefix or suffix sends the results to new columns and leaves the sources alone; with neither, each result overwrites the column it came from. output_data_type casts every result. All the expressions run in one with_columns, so a formula may reference a column that is itself being overwritten and still read its original value. Selecting nothing — an empty columns list, or a data_type no column matches — is a no-op rather than an error.

SQL queries

frame.sql(query) and ff.sql(query, ...) place a SQL Query node and return its output frame. The query uses the Polars SQL dialect and must be a single SELECT or WITH statement.

orders = ff.from_dict(
    {"order_id": [1, 2, 3, 4], "region": ["EU", "US", "EU", "APAC"], "amount": [120.0, 40.0, 900.0, 75.0]}
)
regions = ff.from_dict({"region": ["EU", "US", "APAC"], "target": [800.0, 100.0, 50.0]})

# One frame: it is the table `self`, as in Polars
large = orders.sql("SELECT order_id, region, amount FROM self WHERE amount >= 75")

# Several frames: each keyword names a table
per_region = ff.sql(
    """
    SELECT r.region, SUM(o.amount) AS revenue, r.target
    FROM orders o JOIN regions r ON o.region = r.region
    GROUP BY r.region, r.target
    ORDER BY region
    """,
    orders=large,
    regions=regions,
    description="Revenue per region",
)
result = per_region.collect()
FlowFrame.sql(query: str, *, table_name: str = "self", description: str | None = None) -> FlowFrame

ff.sql(query: str, /, *frames: FlowFrame, description: str | None = None, **tables: FlowFrame) -> FlowFrame
  • frame.sql mirrors polars.LazyFrame.sql: the frame is the table self, or table_name.
  • In ff.sql, positional frames are the tables input_1, input_2, … and each keyword frame is the table of that name, numbered after the positional ones. description is the node label, so it is the one name that cannot be a table. Frames on different graphs are merged onto one. Without frames the query reads no tables and the node starts a new graph: ff.sql("SELECT 1 AS x"). Unlike pl.sql, ff.sql takes its frames as arguments and never reads variables from the calling namespace.
  • The node itself names its inputs input_1, input_2, …, so a named table is stored as a header in front of the query, merged into the query's own WITH when it has one: WITH orders AS (SELECT * FROM input_1), regions AS (SELECT * FROM input_2). The designer shows that header as part of the query. Table names are case-sensitive, as in Polars. The query is trimmed before it is stored, and dedented unless a quoted string in it spans lines.
  • Without description, the node is labelled with the first line of the query.
  • A flow parameter goes in as ${name}, or min_amount.ref in an f-string, and is substituted as text, so a string value needs its SQL quotes: WHERE region = '${region}'. See Flow parameters.
  • An empty query, a table name that is not a plain identifier, a name that is another input's input_<n>, a statement other than SELECT/WITH, and a query that does not resolve against its inputs all raise NativeNodeError when the node is built, and the node is removed again.

String operations

df = df.with_columns([
    ff.col("name").str.to_uppercase().alias("name_upper"),
    ff.col("code").str.slice(0, 3).alias("prefix"),
    ff.col("text").str.contains("pattern").alias("has_pattern"),
])

Conditional logic

df = df.with_columns([
    ff.when(ff.col("price") > 100)
    .then(ff.lit("Premium"))
    .when(ff.col("price") > 50)
    .then(ff.lit("Standard"))
    .otherwise(ff.lit("Budget"))
    .alias("tier"),
])

Date operations

df = df.with_columns([
    ff.col("date").dt.year().alias("year"),
    ff.col("date").dt.month().alias("month"),
    ff.col("date").dt.day().alias("day"),
    ff.col("date").dt.weekday().alias("weekday"),
])

Polars renames apply

The expression namespaces track the pinned Polars version: use dt.weekday() (not day_of_week) and cum_sum() (not cumsum).

List operations

df = df.with_columns([
    ff.col("tags").list.len().alias("tag_count"),
    ff.col("values").list.sum().alias("total"),
    ff.col("items").list.first().alias("first_item"),
])

Polars compatibility

Most Polars Expr methods are available. See the Polars docs for the full method reference; a few methods are renamed or fall back to polars_code nodes — see Expressions.


← Previous: Data Types | Next: Aggregations →