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
polarsexpressions into a JSON-serializable structure.Each
polars.Expris serialized withExpr.meta.serialize(format="binary")and wrapped in a marker dictionary{"__polars_expr__": <base64 string>}, so it can be restored bydeserialize_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
polarsexpressions are encoded.
- b2luigi.contrib.polars.parameters.deserialize_with_polars_value(value)[source]#
Recursively restore the
polarsexpressions in a structure created byserialize_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.Exprobjects reconstructed.
- b2luigi.contrib.polars.parameters.deterministic_polars_hash_function(value) str[source]#
Create a deterministic hash for a value that may contain
polarsexpressions.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 defaulthash_functionof all parameters in this module.Caution
The serialization format of
polarsexpressions is not guaranteed to be stable acrosspolarsversions, so the hash (and with it the output path of a task) may change when you updatepolars.
- class b2luigi.contrib.polars.parameters.PolarsExpressionParameter(*args, **kwargs)[source]#
Bases:
_PolarsSerializationMixin,ParameterParameter for a single
polarsexpression (polars.Expr).Unlike other parameters,
hasheddefaults toTrue, as the serialized expression is not suitable for a file path. The hash is computed withdeterministic_polars_hash_function()unless you provide your ownhash_function. Be aware that the serialization ofpolarsexpressions, and therefore the hash, might not be stable across different versions ofpolars. The same applies toPolarsExpressionListParameterandPolarsExpressionDictParameter.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,ListParameterParameter for a list that may contain
polarsexpressions.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,DictParameterParameter for a dictionary whose (possibly nested) values may be
polarsexpressions.This allows you to store rich configurations with
polars.Exprobjects alongside regular values. LikePolarsExpressionParameter, 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"])