Polars中能否替换计算图源头的LazyFrame?关于序列化延迟执行计划并应用于实际数据的技术咨询
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:
1. Serialize Individual Transformation Logic (Recommended)
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.
2. Hack the Internal Logical Plan (Not Recommended for Production)
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

