1. 项目概述:为什么我们需要另一个DataFrame库?
如果你最近在数据处理圈子里待过,大概率会听到“Polars”这个名字。它不是一个新概念,但正以惊人的速度成为许多数据工程师和分析师工具箱里的新宠。我自己也是从Pandas的深度用户转向Polars的,这个转变的驱动力很简单:当你的数据从GB级别向TB级别迈进,或者当你需要处理实时流数据时,传统的单线程、内存驻留式的Pandas开始显得力不从心。Polars的出现,正是为了解决这个痛点。
Polars是一个用Rust编写的高性能DataFrame库,它通过Apache Arrow作为内存格式,并充分利用了现代CPU的多核并行和向量化计算能力。简单来说,它能让你的数据处理代码跑得更快,尤其是在处理大规模数据集时,速度提升往往是数量级的。它的API设计借鉴了Pandas的易用性,但底层是完全不同的执行引擎。对于数据从业者而言,掌握Polars的常见用法,意味着你能更高效地应对日益增长的数据规模和复杂度挑战。这篇文章,我将结合自己从Pandas迁移到Polars的实战经验,总结那些最常用、最高效的操作模式,帮你快速上手。
2. 核心设计理念与性能基石
在深入具体语法之前,理解Polars为什么快至关重要。这决定了你该如何以“Polars的方式”去思考问题,而不仅仅是把Pandas代码翻译过来。
2.1 惰性求值与查询优化
这是Polars与Pandas最根本的区别之一。Polars提供了两种执行模式:即时执行(Eager)和惰性执行(Lazy)。
- 即时执行模式:类似于Pandas,你输入一个操作,它立刻返回结果。这对于探索性数据分析(EDA)和小型数据集交互非常友好。
- 惰性执行模式(LazyFrame):这是Polars性能的杀手锏。当你创建一个
LazyFrame并对其进行一系列操作(筛选、聚合、连接等)时,Polars并不会立即计算。相反,它会构建一个逻辑执行计划(Logical Plan)。只有当你调用.collect()、.fetch()或.sink_*()方法时,它才会启动优化器。
这个优化器会做几件聪明事:
- 谓词下推(Predicate Pushdown):如果你先
select某些列,再filter行,优化器可能会将过滤条件“下推”到更早的阶段,甚至是在从文件读取数据时就直接应用过滤,大幅减少需要加载和处理的数据量。 - 投影下推(Projection Pushdown):只选择查询最终需要的列,避免将无关列加载到内存中。
- 操作融合:将多个连续的操作合并为一个更高效的操作。
例如,你想从一个大CSV文件中读取数据,过滤出“2023年”的记录,然后按“城市”分组计算“销售额”总和。在惰性模式下,Polars可能会优化为:扫描文件时,只读取“日期”、“城市”、“销售额”三列,并在读取过程中直接过滤掉非2023年的行,最后进行分组聚合。这个优化过程对用户是透明的,你只需要写出逻辑步骤,Polars负责找到最优执行路径。
实操心得:对于生产环境的数据处理管道,务必使用惰性执行模式。它不仅能提升性能,还能让你更清晰地表达数据处理逻辑。可以把
.lazy()看作一个性能开关,习惯性地在数据读取后加上它。
2.2 基于Apache Arrow的列式内存布局
Polars在内存中使用Apache Arrow格式存储数据。这是一种列式存储格式。与Pandas的行式存储(将一行中的所有值连续存放)不同,列式存储将每一列的数据连续存放。
这种布局的好处:
- 极高的缓存利用率:进行聚合运算(如
sum、mean)时,CPU可以连续读取同一列的大量数值,非常适合现代CPU的预取机制,计算速度极快。 - 高效的压缩:同一列的数据类型相同,更容易压缩,节省内存。
- 向量化计算:可以利用SIMD(单指令多数据)指令,让CPU在一个时钟周期内对多个数据执行相同操作,这是Polars许多操作速度远超Pandas的原因。
2.3 并行与无复制操作
Rust语言本身保证了内存安全和无数据竞争的并发。Polars利用这一点,可以安全地将许多操作并行化,例如在多列上应用函数、分组聚合等。同时,许多操作(如切片、增加列)是零复制(Zero-copy)或延迟复制(Lazy Copy)的,避免了不必要的内存分配和数据移动开销。
理解了这些,你就会明白,为什么在Polars中,链式调用(Method Chaining)不仅是风格问题,更是性能最佳实践。因为它允许优化器看到一个完整的操作序列。
3. 从零开始:数据IO与基本操作
让我们从最基础的开始,看看如何用Polars替代你熟悉的Pandas操作。
3.1 数据读取与写入
Polars支持丰富的IO格式,其API设计直观。
import polars as pl # 读取CSV文件 (即时执行) df_eager = pl.read_csv("data.csv") # 读取CSV文件并转为惰性执行模式(推荐用于大数据) df_lazy = pl.scan_csv("data.csv") # 或者 pl.read_csv("data.csv").lazy() # 读取Parquet文件(列式存储,与Polars是天作之合) df_parquet = pl.read_parquet("data.parquet") lazy_parquet = pl.scan_parquet("data.parquet") # 读取JSON df_json = pl.read_json("data.json", format='json') # 或 format='jsonl' # 从Pandas DataFrame转换(桥梁) import pandas as pd pdf = pd.DataFrame({'a': [1,2,3], 'b': ['x', 'y', 'z']}) df_from_pandas = pl.from_pandas(pdf) # 写入数据 df_eager.write_csv("output.csv") df_eager.write_parquet("output.parquet") # 写入Parquet通常是更好的选择 lazy_parquet.sink_parquet("lazy_output.parquet") # 惰性Frame的写入方式注意事项:
scan_csv和scan_parquet直接创建LazyFrame,是处理大文件的起点。- 写入Parquet格式通常比CSV好得多,它压缩率高、读取快,并且能保留数据类型(如日期时间、分类)。
- 从Pandas转换时,注意大DataFrame的内存拷贝开销。对于极大数据集,应优先使用Polars直接读取源文件。
3.2 数据查看与基本属性
df = pl.read_csv("sample_data.csv") # 查看前n行,类似df.head() print(df.head(5)) # 查看形状 print(df.shape) # 查看列名 print(df.columns) # 查看数据类型(Polars中称为dtype) print(df.schema) # 查看统计摘要 print(df.describe())Polars的DataFrame对象是不可变的(immutable)。大多数操作都会返回一个新的DataFrame,这有助于避免副作用,并使代码更易于推理。这与Pandas的inplace=True参数有哲学上的不同。
3.3 列的选择与操作
这是最频繁的操作之一。Polars提供了多种灵活的方式。
# 选择单列(返回一个Series) series_a = df["column_a"] # 选择多列(返回一个DataFrame) df_selected = df[["column_a", "column_b", "column_c"]] # 更Polars风格的写法(支持链式调用) df_selected = df.select(["column_a", "column_b", "column_c"]) # 使用pl.col选择器,功能更强大 df_selected = df.select(pl.col("column_a"), pl.col("column_b") * 2) # 排除某些列 df_without = df.select(pl.exclude("column_to_drop", "another_column")) # 选择所有数值列/字符串列 df_numeric = df.select(pl.col(pl.NUMERIC_DTYPES)) df_string = df.select(pl.col(pl.Utf8)) # 重命名列 df_renamed = df.rename({"old_name": "new_name", "old_name2": "new_name2"})pl.col是Polars表达式的核心,它代表对列的一系列操作,而不是立即计算的值。这种“表达式”可以组合、传递,并在惰性求值中被优化。
4. 数据清洗与转换的实战技巧
数据清洗是数据分析的基石,Polars在这方面提供了强大且高效的工具集。
4.1 过滤行数据:不仅仅是df[df['col'] > 0]
过滤是高频操作。Polars的过滤语法直观且强大。
# 基础过滤 df_filtered = df.filter(pl.col("age") > 18) # 多条件组合 (使用 &, |, ~ 代替 and, or, not) df_complex = df.filter( (pl.col("age") > 18) & (pl.col("city").is_in(["北京", "上海"])) & (~pl.col("name").str.contains("测试")) ) # 过滤空值 df_non_null = df.filter(pl.col("salary").is_not_null()) # 或者直接删除包含空值的行(谨慎使用,可能删除大量数据) df_dropped = df.drop_nulls() # 根据热搜词“指定数据开头过滤”:过滤出某列以特定字符串开头的行 # 例如,过滤出`user_id`以‘UA’开头的记录 df_startswith = df.filter(pl.col("user_id").str.starts_with("UA")) # 同理,还有 .str.ends_with() 和 .str.contains()实操心得:在惰性模式下,过滤条件会尽可能地被“下推”到数据源。这意味着如果你从Parquet文件读取并立即过滤,Polars可能只读取满足条件的行所在的数据页,而不是整个文件,这对性能提升是巨大的。务必在
scan之后尽早进行filter。
4.2 处理缺失值与数据类型转换
# 填充空值 df_filled = df.with_columns( pl.col("salary").fill_null(0), # 用0填充 pl.col("name").fill_null("Unknown"), # 用字符串填充 pl.col("date").forward_fill(), # 用前一个有效值向前填充 ) # 更复杂的填充策略:按分组填充均值 df_group_fill = df.with_columns( pl.col("salary").fill_null(pl.col("salary").mean().over("department")) ) # 数据类型转换 df_converted = df.with_columns( pl.col("price").cast(pl.Float64), # 转换为浮点 pl.col("timestamp_str").str.strptime(pl.Datetime, format="%Y-%m-%d %H:%M:%S"), # 字符串转日期时间 pl.col("int_col").cast(pl.Utf8), # 整型转字符串 )4.3 创建新列与列操作
with_columns方法是Polars中新增列或修改现有列的主力,它返回一个新的DataFrame,包含所有原有列以及新增或修改的列。
# 创建简单的新列 df_new = df.with_columns( (pl.col("price") * pl.col("quantity")).alias("revenue"), # 计算收入 (pl.col("date").dt.year()).alias("year") # 提取年份 ) # 使用条件逻辑创建列(类似np.where或pandas的df.apply) df_with_logic = df.with_columns( pl.when(pl.col("revenue") > 1000) .then("High") .when(pl.col("revenue") > 500) .then("Medium") .otherwise("Low") .alias("revenue_tier") ) # 对字符串列进行操作 df_string_ops = df.with_columns( pl.col("email").str.split("@").list.get(1).alias("domain"), # 提取邮箱域名 pl.col("name").str.to_uppercase().alias("name_upper"), pl.col("description").str.replace_all(r"\s+", " ").alias("desc_clean") # 替换多余空格 ) # 对列表列进行操作(如果某列是List类型) df_list_ops = df.with_columns( pl.col("tags").list.lengths().alias("num_tags"), # 列表长度 pl.col("scores").list.mean().alias("avg_score") # 列表内均值 )with_columns的强大之处在于,它接受多个表达式,并且这些表达式可以引用在同一调用中刚刚创建的其他列。Polars的优化器会处理这些依赖关系。
5. 分组、聚合与窗口函数
分组聚合是数据分析的核心。Polars的分组聚合性能极其出色,得益于其列式存储和并行计算。
5.1 基础分组聚合
# 单维度分组,单指标聚合 df_grouped = df.group_by("department").agg( pl.col("salary").mean().alias("avg_salary"), pl.col("salary").sum().alias("total_salary"), pl.col("employee_id").count().alias("headcount") ) # 多维度分组 df_multi_group = df.group_by("year", "quarter", "region").agg( pl.col("sales").sum().alias("total_sales"), pl.col("profit").mean().alias("avg_profit") ) # 多个聚合函数应用于同一列 df_multi_agg = df.group_by("category").agg( pl.col("price").min().alias("min_price"), pl.col("price").max().alias("max_price"), pl.col("price").mean().alias("mean_price"), pl.col("price").std().alias("std_price") )5.2 高级聚合与表达式
Polars的聚合表达式非常灵活,你可以在聚合内部进行复杂的计算。
# 聚合时进行过滤:只聚合满足条件的记录 # 例如,计算每个部门“高薪”(>50000)员工的平均工资 df_cond_agg = df.group_by("department").agg( pl.col("salary").filter(pl.col("salary") > 50000).mean().alias("avg_high_salary") ) # 聚合后排序 df_sorted_agg = ( df.group_by("department") .agg(pl.col("salary").sum().alias("total_salary")) .sort("total_salary", descending=True) # 按聚合结果降序排序 )5.3 窗口函数:不减少行数的“分组”
窗口函数允许你在每一行上执行计算,同时参考一个与当前行相关的行“窗口”。这是进行排名、移动平均、累计求和等操作的利器。
# 排名:每个部门内按工资排名 df_rank = df.with_columns( pl.col("salary").rank(method="dense").over("department").alias("dept_salary_rank") ) # 移动平均:计算每个产品最近3天的销售额移动平均(假设数据已按日期排序) df_ma = df.sort("date").with_columns( pl.col("daily_sales").rolling_mean(window_size=3, min_periods=1).over("product_id").alias("sales_ma_3d") ) # 累计求和:计算每个用户订单金额的累计和 df_cumsum = df.sort("order_date").with_columns( pl.col("amount").cum_sum().over("user_id").alias("cumulative_amount") ) # 组内偏移:获取每个用户上一次订单的金额(lag) df_lag = df.sort("order_date").with_columns( pl.col("amount").shift(1).over("user_id").alias("prev_order_amount") )注意事项:窗口函数中的
.over(“group_col”)是关键。它定义了窗口的划分范围。在计算移动窗口统计量(如rolling_mean)时,必须确保数据在组内已按时间顺序排序,否则结果毫无意义。我建议在应用窗口函数前,先进行.sort([“group_col”, “time_col”])操作。
6. 表连接与数据合并
将多个数据集合并是常见任务。Polars支持多种连接类型,语法清晰。
6.1 多种连接方式
df_left = pl.DataFrame({ "key": ["A", "B", "C", "D"], "value_left": [1, 2, 3, 4] }) df_right = pl.DataFrame({ "key": ["B", "C", "D", "E"], "value_right": [5, 6, 7, 8] }) # 内连接 (inner join):只保留两个表都有的key df_inner = df_left.join(df_right, on="key", how="inner") # 左连接 (left join):保留左表所有行,右表匹配不上则为null df_left_join = df_left.join(df_right, on="key", how="left") # 全外连接 (outer join):保留所有行,缺失处为null df_outer = df_left.join(df_right, on="key", how="outer") # 半连接 (semi join):只保留左表中那些在右表有关联键的行,不添加右表的列 df_semi = df_left.join(df_right, on="key", how="semi") # 反连接 (anti join):只保留左表中那些在右表没有关联键的行 df_anti = df_left.join(df_right, on="key", how="anti") # 使用多个键进行连接 df_multi_key = df_left.join(df_right, on=["key1", "key2"], how="inner")6.2 连接的性能考量与重复列名
# 当连接键在两个表中列名不同时 df_left.join(df_right, left_on="left_key", right_on="right_key", how="inner") # 处理连接后的重复列名(非连接键列名相同) df_left = pl.DataFrame({"a": [1,2], "b": [3,4]}) df_right = pl.DataFrame({"a": [1,2], "b": [5,6]}) # 列`b`重复 result = df_left.join(df_right, on="a", how="inner", suffix="_right") # 结果中,列名会变为 `b` 和 `b_right`实操心得:对于超大型表的连接,性能是关键。
- 惰性连接:在
LazyFrame上使用.join(),优化器可能会将过滤条件下推到连接之前,或者选择更高效的连接算法(如哈希连接)。- 广播连接:如果右表非常小,Polars可能会自动采用“广播连接”,将小表复制到所有工作线程,这通常很快。你可以通过设置
how=“inner_coalesce”等策略给予提示。- 连接前过滤:尽可能在连接前使用
filter减少每个表的数据量,这是提升连接速度最有效的方法之一。
7. 惰性执行(LazyFrame)的深入应用
如前所述,惰性执行是处理大数据时的首选。让我们看看如何构建和优化一个完整的惰性查询。
7.1 构建一个完整的惰性查询管道
# 1. 从源创建LazyFrame lazy_df = pl.scan_parquet("large_data.parquet") # 2. 构建查询计划(此时没有实际计算) query = (lazy_df .filter(pl.col("date").dt.year() == 2023) # 尽早过滤 .filter(pl.col("status") == "active") .select(["user_id", "department", "amount", "date"]) # 只选择需要的列 .with_columns( (pl.col("amount") * 1.1).alias("amount_with_tax"), # 计算新列 pl.col("date").dt.month().alias("month") ) .group_by("department", "month") .agg( pl.col("amount_with_tax").sum().alias("total_revenue"), pl.col("user_id").n_unique().alias("unique_users") ) .sort("total_revenue", descending=True) ) # 3. 查看优化前的逻辑计划(用于调试) print(query.explain()) # 4. 查看优化后的物理计划(更接近实际执行) print(query.explain(optimized=True)) # 5. 触发计算并获取结果 result_df = query.collect() # 将所有结果拉取到内存 # 或者,如果只想查看一部分 sample_result = query.fetch(n_rows=1000) # 适合预览,可能不是最终结果的随机样本 # 或者,直接写入磁盘(对于非常大的结果集) query.sink_parquet("aggregated_results.parquet").explain()是你的好朋友。通过查看计划,你可以了解Polars将如何执行你的查询,有时可以发现优化空间(比如过滤条件的位置是否最优)。
7.2 惰性模式下的常见优化技巧
- 谓词下推:确保过滤操作(
filter)尽可能早地出现在链中,最好紧接在scan之后。这样,数据源连接器(如scan_parquet)可能直接在读取时应用过滤。 - 投影下推:尽早使用
select明确指定你需要的列,避免将整行数据(尤其是包含大文本字段的列)带入后续计算。 - 避免在惰性帧上使用
.to_pandas()或.collect()中间结果:这会打断优化计划,强制进行物化计算。应保持完整的操作链,最后再collect。 - 使用
.sink_*进行流式输出:对于最终输出到文件的操作,使用sink_parquet或sink_ipc可以让Polars以流式方式写入,避免在内存中物化整个结果集。
8. 性能调优与常见陷阱
即使使用了Polars,不当的使用方式也可能导致性能不佳。以下是一些关键的性能调优点和常见“坑”。
8.1 选择正确的数据类型
Polars的数据类型(dtype)直接影响内存占用和计算速度。
- 数值类型:使用能满足需求的最小类型。例如,如果数值范围在0-255,用
pl.UInt8而非pl.Int64,内存占用减少为1/8。 - 字符串类型:
pl.Utf8是通用字符串。如果字符串是分类变量且基数(唯一值数量)不大,考虑转换为pl.Categorical类型,可以显著提升分组和过滤速度,并减少内存。 - 日期时间:使用
pl.Date、pl.Datetime、pl.Duration等专门类型,而不是字符串,以便利用时间序列优化函数。
# 优化数据类型示例 df_optimized = df.with_columns( pl.col("category").cast(pl.Categorical), # 分类列转换 pl.col("small_int").cast(pl.UInt8), pl.col("timestamp_str").str.strptime(pl.Datetime(time_unit="us")) # 明确时间单位 )8.2 避免行级迭代,使用向量化操作
这是从Pandas迁移过来最容易犯的错误。在Pandas中,df.apply()或循环有时难以避免,但在Polars中,这将是性能灾难。
# **错误示范** (极慢!) result = [] for row in df.iter_rows(): # 或 df.to_dicts() # 对每一行进行复杂计算... pass # **正确示范** (向量化,极快!) # 使用 when().then().otherwise() # 使用 pl.col().map_elements() 仅作为最后手段,且确保提供的函数是经过优化的(如numpy函数) # 绝大多数逻辑都可以用内置表达式完成 df = df.with_columns( pl.when(pl.col("x") > pl.col("y")) .then(pl.col("x") - pl.col("y")) .otherwise(pl.col("y") - pl.col("x")) .alias("diff") )如果确实需要应用一个复杂的自定义函数,并且无法用内置表达式实现,可以考虑:
- 使用
pl.col().map_elements(function, return_dtype=...),但这是最后的选择,因为它会强制将数据传递到Python端,损失性能。 - 使用Polars的
Struct类型或list.eval来处理更复杂的行内逻辑。
8.3 内存管理
- 流式处理:对于远超内存的数据,使用
scan_*创建LazyFrame,并通过.sink_parquet()流式输出,或使用.collect(streaming=True)进行流式收集(需要配置)。 - 分块处理:如果必须使用即时执行模式,可以考虑手动将数据分块处理。
- 监控内存:使用
df.estimated_size(“mb”)来查看DataFrame的预估内存占用。
8.4 序列化与反序列化
在分布式计算或缓存中间结果时,序列化格式很重要。
- Parquet:是磁盘存储和交换的最佳选择,压缩率高,Polars读写极快。
- IPC/Feather格式(.arrow, .feather):这是Apache Arrow的二进制格式,序列化和反序列化速度最快,适合在内存或高速存储中暂存数据。使用
df.write_ipc()和pl.read_ipc()。
# 快速缓存中间结果到本地 df.write_ipc(“intermediate.arrow”) df_fast_load = pl.read_ipc(“intermediate.arrow”)9. 与生态系统的集成
Polars不是孤岛,它需要与现有工具链协同工作。
9.1 与Pandas互操作
虽然鼓励直接使用Polars,但有时不得不与依赖Pandas的库交互。
# Polars -> Pandas pandas_df = df.to_pandas(use_pyarrow_extension_array=True) # 使用PyArrow扩展数组,转换更快 # Pandas -> Polars polars_df = pl.from_pandas(pandas_df)注意:
to_pandas()会将所有数据从Arrow内存格式复制到Pandas的NumPy格式中。对于大型DataFrame,这是一个昂贵操作,可能导致内存峰值。仅在必要时使用。
9.2 与SQL交互
Polars内置了一个小型SQL引擎,可以用SQL查询DataFrame或LazyFrame。
# 注册DataFrame/LazyFrame为一个临时表 df = pl.DataFrame({"a": [1,2,3], "b": [4,5,6]}) ctx = pl.SQLContext(my_table=df) # 注册df为`my_table` # 执行SQL查询 result = ctx.execute("SELECT a, b*2 as b_double FROM my_table WHERE a > 1") print(result.collect())这对于熟悉SQL的团队快速上手或执行一些复杂的多表连接查询非常方便。但要注意,为了获得最佳性能,特别是利用惰性求值和优化器,原生Polars表达式API仍然是首选。
9.3 可视化
Polars DataFrame可以无缝转换为Pandas DataFrame,从而利用成熟的Matplotlib, Seaborn, Plotly等可视化库。对于简单的预览,Polars也提供了.plot()方法(需要安装pyarrow和matplotlib)。
10. 实战案例:一个端到端的用户行为分析片段
让我们用一个模拟的电商用户行为日志,串联起多个常见操作。
假设我们有一个user_logs.parquet文件,包含字段:user_id,session_id,event_time,event_type(‘click’, ‘view’, ‘purchase’),product_id,amount。
目标:计算2023年第二季度,每个用户的购买总金额、购买次数、以及最后一次购买前7天内的点击事件总数。
# 使用惰性执行从Parquet读取 lazy_logs = pl.scan_parquet("user_logs.parquet") analysis_result = ( lazy_logs # 1. 过滤时间和事件类型,选择所需列 .filter( (pl.col("event_time").dt.year() == 2023) & (pl.col("event_time").dt.quarter() == 2) ) .select(["user_id", "event_time", "event_type", "amount"]) # 2. 分离购买事件和点击事件 # 我们创建两个“虚拟列”,一个用于购买聚合,一个用于点击窗口计算 .with_columns( # 标记是否为购买事件,并携带金额 pl.when(pl.col("event_type") == "purchase") .then(pl.col("amount")) .otherwise(0) .alias("purchase_amount"), # 标记是否为点击事件 (pl.col("event_type") == "click").cast(pl.UInt8).alias("is_click") ) # 3. 按用户分组,进行聚合 .group_by("user_id") .agg( # 购买总金额 pl.col("purchase_amount").sum().alias("total_purchase_amount"), # 购买次数(金额非0的次数) pl.col("purchase_amount").count().alias("total_purchase_count"), # 获取每个用户最后一次购买的时间 pl.col("event_time") .filter(pl.col("event_type") == "purchase") .max() .alias("last_purchase_time") ) # 4. 将聚合结果与原始日志再次连接,计算窗口点击量 # 这里需要将聚合结果(每个用户一行)与原始日志(每个事件一行)连接 # 为了演示,我们假设数据量可以接受,先collect聚合结果。对于超大数据,有更高级的窗口函数写法。 ).collect() # 注意:上面的查询只完成了聚合。要计算“最后一次购买前7天的点击”, # 更高效的写法是在一个复杂的窗口函数中完成,但为了清晰,我们分步演示。 # 在实际生产中,应尝试在一个查询内用高级窗口函数完成。 # 假设 analysis_result 不大,我们可以进行二次连接计算 lazy_logs_clicks = lazy_logs.filter(pl.col("event_type") == "click") # 将最后一次购买时间广播回去,然后过滤计算 final_result = ( lazy_logs_clicks .join(analysis_result.lazy(), on="user_id", how="inner") .filter( (pl.col("event_time") > pl.col("last_purchase_time") - pl.duration(days=7)) & (pl.col("event_time") < pl.col("last_purchase_time")) ) .group_by("user_id") .agg( pl.col("event_time").count().alias("clicks_7d_before_last_purchase") ) .join(analysis_result.lazy(), on="user_id", how="left") .select(["user_id", "total_purchase_amount", "total_purchase_count", "clicks_7d_before_last_purchase"]) .collect() ) print(final_result)这个案例展示了过滤、条件列创建、分组聚合、多次连接和条件过滤的组合。在真实场景中,对于最后一步的窗口计算,可以研究使用pl.col(“event_time”).filter(…).count().over(“user_id”)配合复杂的窗口定义来尝试一次性完成,避免collect中间结果。这需要根据数据分布和大小进行权衡和测试。
迁移到Polars是一个思维转换的过程,从“行式迭代”转向“列式向量化”和“声明式查询优化”。开始时可能会觉得有些约束,但一旦习惯,其带来的性能提升和代码清晰度会让你觉得物超所值。从今天开始,尝试在你的下一个数据任务中使用Polars,先从替换一个Pandas的read_csv和groupby开始,你会立刻感受到不同。