如何用Python Polars实现无demean的分组自定义相关系数计算?
在Python Polars中实现自定义无demean相关系数计算
下面给你三种可行的实现方式,对应你提到的map_elements、map_batches(通过group_by.apply)和自定义表达式命名空间的用法:
方法1:在agg中使用map_elements
直接在分组聚合时,把需要的列打包成struct,通过map_elements调用自定义函数:
不加权场景
df.group_by(["UNIVERSE", "datetime"]).agg( corr_xy = pl.struct(["ratio_wsm", "y1d_nn_r"]).map_elements( lambda s: rcor(s["ratio_wsm"].to_numpy(), s["y1d_nn_r"].to_numpy()), return_dtype=pl.Float64 ) )
加权场景(假设有权重列weight)
df.group_by(["UNIVERSE", "datetime"]).agg( corr_xy = pl.struct(["ratio_wsm", "y1d_nn_r", "weight"]).map_elements( lambda s: rcor(s["ratio_wsm"].to_numpy(), s["y1d_nn_r"].to_numpy(), s["weight"].to_numpy()), return_dtype=pl.Float64 ) )
说明:pl.struct将分组内的目标列打包成一个结构体,map_elements对每个分组的结构体执行自定义逻辑,最后指定返回数据类型。
方法2:使用group_by.apply(基于map_batches)
通过group_by.apply接收每个分组的完整DataFrame,更灵活处理复杂逻辑:
def compute_rcor_per_group(group_df: pl.DataFrame) -> pl.DataFrame: # 提取分组内的列数据 x = group_df["ratio_wsm"].to_numpy() y = group_df["y1d_nn_r"].to_numpy() # 处理权重可选的情况 w = group_df["weight"].to_numpy() if "weight" in group_df.columns else None # 计算自定义相关系数 corr_val = rcor(x, y, w) # 返回包含分组键和结果的小DataFrame return pl.DataFrame({ "UNIVERSE": [group_df["UNIVERSE"][0]], "datetime": [group_df["datetime"][0]], "corr_xy": [corr_val] }) # 执行分组计算 result = df.group_by(["UNIVERSE", "datetime"], maintain_order=True).apply(compute_rcor_per_group)
说明:apply会遍历每个分组的DataFrame,你可以在自定义函数中对分组数据做任意操作,最后返回的小DataFrame会被Polars自动合并成最终结果。
方法3:注册自定义表达式命名空间
如果想和Polars原生函数(如pl.corr)一样用表达式风格调用,可以注册自己的命名空间,提升代码复用性和可读性:
第一步:注册命名空间
@pl.api.register_expr_namespace("custom") class CustomExpr: def __init__(self, expr: pl.Expr): self._expr = expr def rcor(self, other: pl.Expr, weight: pl.Expr = None) -> pl.Expr: if weight is not None: return pl.struct([self._expr, other, weight]).map_elements( lambda s: rcor(s[0].to_numpy(), s[1].to_numpy(), s[2].to_numpy()), return_dtype=pl.Float64 ) else: return pl.struct([self._expr, other]).map_elements( lambda s: rcor(s[0].to_numpy(), s[1].to_numpy()), return_dtype=pl.Float64 )
第二步:像原生函数一样调用
# 不加权场景 result = df.group_by(["UNIVERSE", "datetime"]).agg( corr_xy = pl.col("ratio_wsm").custom.rcor(pl.col("y1d_nn_r")) ) # 加权场景 result = df.group_by(["UNIVERSE", "datetime"]).agg( corr_xy = pl.col("ratio_wsm").custom.rcor(pl.col("y1d_nn_r"), pl.col("weight")) )
性能优化建议
你的自定义函数目前用numpy实现,其实可以改成纯Polars原生操作,避免numpy和Polars之间的类型转换开销,性能会更好:
def rcor_polars(x: pl.Series, y: pl.Series, w: pl.Series = None) -> float: if w is not None: sxx = (w * x * x).sum() syy = (w * y * y).sum() sxy = (w * x * y).sum() else: sxx = (x * x).sum() syy = (y * y).sum() sxy = (x * y).sum() return (sxy / pl.sqrt(sxx * syy)).item()
使用时只需要把之前的rcor替换成rcor_polars,并且不用转numpy,直接传Series即可,比如方法1中的lambda可以改成:
lambda s: rcor_polars(s["ratio_wsm"], s["y1d_nn_r"])
内容的提问来源于stack exchange,提问作者LeoKnuth
相关产品推荐
相关产品推荐

