polars#

The b2luigi.contrib.polars module provides parameters that accept polars expressions (polars.Expr), either on their own or inside (nested) lists and dictionaries. This lets you, e.g., define selections or derived columns as task parameters:

pip3 install "b2luigi[polars]"
import polars as pl
import b2luigi
from b2luigi.contrib.polars import PolarsExpressionParameter


class ApplySelection(b2luigi.Task):
    selection = PolarsExpressionParameter()

    def output(self):
        yield self.add_to_output("selected.parquet")

    def run(self):
        df = pl.read_parquet("input.parquet")
        df.filter(self.selection).write_parquet(self.get_output_file_name("selected.parquet"))


if __name__ == "__main__":
    b2luigi.process(ApplySelection(selection=pl.col("x") > 0))

The expressions are serialized with polars’ binary format, so they can be passed to batch jobs. As serialized expressions are not readable in a file path, all parameters in this package are hashed by default (see the hashed flag described in wrap_parameter).

Caution

The binary serialization format of polars expressions is not guaranteed to be stable across polars versions. Updating polars may therefore change the hashes, and with them the output paths, of tasks with polars expression parameters.

API Reference#

Parameters for handling polars expressions (polars.Expr).

To learn more about polars expressions, visit the polars documentation.

b2luigi.contrib.polars.parameters.POLARS_EXPR_KEY = '__polars_expr__'#

Key of the marker dictionary that wraps a serialized polars.Expr.

b2luigi.contrib.polars.parameters.serialize_with_polars_value(value)[source]#

Recursively convert a value that may be or may contain polars expressions into a JSON-serializable structure.

Each polars.Expr is serialized with Expr.meta.serialize(format="binary") and wrapped in a marker dictionary {"__polars_expr__": <base64 string>}, so it can be restored by deserialize_with_polars_value(). Mappings, lists and tuples are traversed recursively, all other values are returned unchanged.

Parameters:

value – Any value, e.g. a polars.Expr, a mapping, list or tuple, or a scalar.

Returns:

A JSON-serializable structure in which all polars expressions are encoded.

b2luigi.contrib.polars.parameters.deserialize_with_polars_value(value)[source]#

Recursively restore the polars expressions in a structure created by serialize_with_polars_value().

As JSON has no tuples, sequences are always returned as lists. The parameters in this module freeze them to tuples anyway when a task is instantiated.

Parameters:

value – A structure previously produced by serialize_with_polars_value().

Returns:

The original structure with all polars.Expr objects reconstructed.

b2luigi.contrib.polars.parameters.deterministic_polars_hash_function(value) → str[source]#

Create a deterministic hash for a value that may contain polars expressions.

The value is serialized with serialize_with_polars_value() and the first 16 characters of the SHA-256 digest of the (key-sorted) JSON representation are returned. This is the default hash_function of all parameters in this module.

Caution

The serialization format of polars expressions is not guaranteed to be stable across polars versions, so the hash (and with it the output path of a task) may change when you update polars.

class b2luigi.contrib.polars.parameters.PolarsExpressionParameter(*args, **kwargs)[source]#

Bases: _PolarsSerializationMixin, Parameter

Parameter for a single polars expression (polars.Expr).

Unlike other parameters, hashed defaults to True, as the serialized expression is not suitable for a file path. The hash is computed with deterministic_polars_hash_function() unless you provide your own hash_function. Be aware that the serialization of polars expressions, and therefore the hash, might not be stable across different versions of polars. The same applies to PolarsExpressionListParameter and PolarsExpressionDictParameter.

Example

import polars as pl
import b2luigi
from b2luigi.contrib.polars import PolarsExpressionParameter

class MyTask(b2luigi.Task):
    expr = PolarsExpressionParameter(default=pl.col("x") * 2)

    def run(self):
        df = pl.DataFrame({"x": [1, 2, 3]})
        result = df.select(self.expr.alias("double_x"))
        print(result)
class b2luigi.contrib.polars.parameters.PolarsExpressionListParameter(*args, **kwargs)[source]#

Bases: _PolarsSerializationMixin, ListParameter

Parameter for a list that may contain polars expressions.

Like PolarsExpressionParameter, it is hashed by default, and the same caveats on hashing apply.

Example

import polars as pl
import b2luigi
from b2luigi.contrib.polars import PolarsExpressionListParameter

class MyTask(b2luigi.Task):
    filters = PolarsExpressionListParameter(
        default=[
            pl.col("x") > 0,
            pl.col("y").is_not_null(),
        ],
    )

    def run(self):
        df = pl.DataFrame({"x": [1, -2, 3], "y": [None, 5, 6]})
        for filter_expr in self.filters:
            df = df.filter(filter_expr)
        print(df)
class b2luigi.contrib.polars.parameters.PolarsExpressionDictParameter(*args, **kwargs)[source]#

Bases: _PolarsSerializationMixin, DictParameter

Parameter for a dictionary whose (possibly nested) values may be polars expressions.

This allows you to store rich configurations with polars.Expr objects alongside regular values. Like PolarsExpressionParameter, it is hashed by default, and the same caveats on hashing apply.

Example

import polars as pl
import b2luigi
from b2luigi.contrib.polars import PolarsExpressionDictParameter

class MyTask(b2luigi.Task):
    config = PolarsExpressionDictParameter(
        default={
            "compute": pl.col("x") * 3,
            "filter": pl.col("y").is_between(0, 10),
            "settings": {"threshold": 0.95},
        },
    )

    def run(self):
        df = pl.DataFrame({"x": [1, 2, 3], "y": [5, 15, 7]})
        result = df.filter(self.config["filter"]).select(self.config["compute"].alias("triple_x"))
        print(result)
        print(self.config["settings"]["threshold"])