You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Polars中能否替换计算图源头的LazyFrame?关于序列化延迟执行计划并应用于实际数据的技术咨询

Solution for Serializing Polars Lazy Execution Plans and Reapplying to New Data

Great question—this is a super useful use case for edge and lab data scenarios where you want to decouple pipeline logic from deployment. Let's walk through what's possible with Polars today, the underlying constraints, and workarounds to achieve your goal.

Why Your Current Approach Fails

Polars' LazyFrame.serialize() captures the entire query context, including the original source (in your case, the empty LazyFrame). When you deserialize it, the plan is still tied to that empty source—so trying to "attach" it to new data doesn't work because the plan's root node is fixed to the original empty scan.

Feasible Workarounds

Here are two reliable ways to implement your desired pipeline serialization:

Instead of serializing the entire LazyFrame, serialize just the transformation steps (expressions and operation types). You can store these steps in a structured format (like a list of dictionaries) and use Polars' native Expr.serialize() to save expressions safely.

Example implementation:

import polars as pl
import io
import msgpack

# Define your pipeline as a list of operation specs
pipeline_spec = [
    {
        "op": "with_columns",
        "exprs": [pl.col('a') + 1]
    },
    {
        "op": "filter",
        "expr": pl.col('a') > 2
    }
]

# Serialize the pipeline: convert Exprs to bytes
def serialize_pipeline(spec):
    serialized = []
    for step in spec:
        if step["op"] == "with_columns":
            serialized_exprs = [e.serialize() for e in step["exprs"]]
            serialized.append({"op": "with_columns", "exprs": serialized_exprs})
        elif step["op"] == "filter":
            serialized_expr = step["expr"].serialize()
            serialized.append({"op": "filter", "expr": serialized_expr})
    return msgpack.packb(serialized)

# Deserialize and apply to new data
def deserialize_and_apply(serialized_data, target_lf):
    deserialized = msgpack.unpackb(serialized_data)
    current_lf = target_lf
    for step in deserialized:
        if step["op"] == "with_columns":
            exprs = [pl.Expr.deserialize(io.BytesIO(e)) for e in step["exprs"]]
            current_lf = current_lf.with_columns(exprs)
        elif step["op"] == "filter":
            expr = pl.Expr.deserialize(io.BytesIO(step["expr"]))
            current_lf = current_lf.filter(expr)
    return current_lf

# Usage
serialized = serialize_pipeline(pipeline_spec)
actual_data = pl.LazyFrame({'a': [1,2,3]})
processed_data = deserialize_and_apply(serialized, actual_data).collect()
print(processed_data)
# Output:
# shape: (1, 1)
# ┌─────┐
# │ a   │
# │ --- │
# │ i64 │
# ╞═════╡
# │ 4   │
# └─────┘

This approach is stable, uses Polars' public API, and lets you fully decouple the pipeline logic from the source data.

If you need a more direct way to modify the existing query plan, you can access Polars' internal LogicalPlan structure to replace the root source node. Note: This relies on unstable internal APIs that may break in future Polars versions.

Example:

import polars as pl
from polars.planner import LogicalPlan

def replace_plan_source(original_plan: LogicalPlan, new_source_plan: LogicalPlan) -> LogicalPlan:
    """Recursively replace EmptyScan nodes with a new source plan."""
    if isinstance(original_plan, pl.planner.LogicalPlan.EmptyScan):
        return new_source_plan
    # Handle common operation types (Filter, Projection, etc.)
    elif isinstance(original_plan, pl.planner.LogicalPlan.Filter):
        updated_input = replace_plan_source(original_plan.input, new_source_plan)
        return pl.planner.LogicalPlan.Filter(predicate=original_plan.predicate, input=updated_input)
    elif isinstance(original_plan, pl.planner.LogicalPlan.Projection):
        updated_input = replace_plan_source(original_plan.input, new_source_plan)
        return pl.planner.LogicalPlan.Projection(exprs=original_plan.exprs, input=updated_input)
    # Add more cases for other operation types as needed
    else:
        return original_plan

# Define original empty-based plan
empty_lf = pl.LazyFrame()
step1_lf = empty_lf.with_columns(pl.col('a') + 1)
step2_lf = step1_lf.filter(pl.col('a') > 2)

# Create new source data's plan
actual_data = pl.LazyFrame({'a': [1,2,3]})
new_source_plan = actual_data._plan

# Replace the source and build a new LazyFrame
updated_plan = replace_plan_source(step2_lf._plan, new_source_plan)
new_lf = pl.LazyFrame._from_plan(updated_plan)
processed_data = new_lf.collect()
print(processed_data)

Use this only if you can tolerate breaking changes when Polars updates.

Underlying Mechanisms Limiting Native Support

Polars' LazyFrame is designed as a combination of a query plan and its source context (schema, scan settings, etc.). The serialization format includes all this context, so there's no built-in way to "swap" the source without modifying the plan structure. The Polars team has discussed adding support for standalone query plans, but this isn't available as of now.

Can You Replace the Source of a LazyFrame's Computation Graph?

There's no public API for this, but as shown above, you can modify the internal LogicalPlan (with caveats). The recommended approach is to separate your transformation logic from the source data entirely, as outlined in the first workaround.

内容的提问来源于stack exchange,提问作者wterry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.27 10:57:29