理解Polars全新的Categorical类型

Polars的底层实现正在快速演进。去年我们构建了新的流式执行引擎,现在正努力为Polars API打造高性能的分布式引擎。这些工作使此前针对Categorical数据类型的设计决策面临压力。

Polars的流式引擎会并行处理称为morsel的小数据块。过去,每个数据块都带有自己的映射表,因此必须不断同步和重新编码。另一个选择是全局StringCache,但它会在流式流水线中引入锁和停顿。而在分布式架构下,这种以进程为单位的“全局”方案实际上并不真正全局,也正是我们希望避免的设计。

因此,我们重新实现了Categorical数据类型的工作方式(#23016)。本文解释这次改动的原因,以及我们如何解决问题。

最初的实现

Categorical用于表示取值种类有限的字符串列。当唯一值数量远小于列的总条目数时,这种类型很有用:它利用这一特征,以比普通字符串更高效的方式存储数据。具体来说,底层使用较小的数值类型存储,并保留一份映射,将“物理类型”的值(例如8位无符号整数)映射回原始字符串。

先创建一个DataFrame,其中包含Categorical类型的Series。这里使用Polars 1.31.0,它是此次Categorical重构之前的版本。

bears = (
    pl.DataFrame({
        "bears": ["Polar", "Grizzly", "Sun", "Polar"]
    })
    .cast(
        pl.Categorical()
    )
)
print(
    bears.with_columns(
        pl.col("bears").to_physical().alias("physical")
    )
)

# Output
shape: (4, 2)
┌─────────┬──────────┐
│ bears   ┆ physical │
│ ---     ┆ ---      │
│ cat     ┆ u32      │
╞═════════╪══════════╡
│ Polar   ┆ 0        │
│ Grizzly ┆ 1        │
│ Sun     ┆ 2        │
│ Polar   ┆ 0        │
└─────────┴──────────┘

在这个DataFrame中,bears Series的数据类型是Categorical。它有自己的映射,关联物理类型uint32的值与对应字符串。例如,Polar熊被编码为0。

创建另一个含有Categorical列的DataFrame时,会创建新的映射,并将它局部存储在新的Series中。

other_bears = (
    pl.DataFrame({
        "bears": ["Brown", "Spectacled", "Polar", "Sloth"]
    })
    .cast(
        pl.Categorical()
    )
)

print(
    other_bears.with_columns(
        pl.col("bears").to_physical().alias("physical")
    )
)

# Output
┌────────────┬──────────┐
│ bears      ┆ physical │
│ ---        ┆ ---      │
│ cat        ┆ u32      │
╞════════════╪══════════╡
│ Brown      ┆ 0        │
│ Spectacled ┆ 1        │
│ Polar      ┆ 2        │
│ Sloth      ┆ 3        │
└────────────┴──────────┘

这里可以看到,Polar熊的物理值是2。

尝试合并这两个DataFrame时,会遇到以下情况:

pl.concat([bears, other_bears]).with_columns(
    pl.col("bears").to_physical().alias("physical")
)
<sys>:0: CategoricalRemappingWarning: Local categoricals have different encodings, 
expensive re-encoding is done to perform this merge operation. 
Consider using a StringCache or an Enum type if the categories are known in advance

# Output
shape: (8, 2)
┌────────────┬──────────┐
│ bears      ┆ physical │
│ ---        ┆ ---      │
│ cat        ┆ u32      │
╞════════════╪══════════╡
│ Polar      ┆ 0        │
│ Grizzly    ┆ 1        │
│ Sun        ┆ 2        │
│ Polar      ┆ 0        │
│ Brown      ┆ 3        │
│ Spectacled ┆ 4        │
│ Polar      ┆ 0        │
│ Sloth      ┆ 5        │
└────────────┴──────────┘

为了合并这些Categorical,需要同步它们的映射,确保新合并的映射中每个类别都有正确表示。由于重新编码的开销很大,我们发出警告,建议使用StringCache:它是一份全局映射,所有DataFrame都能使用。

使用StringCache后,不同DataFrame中的物理值会保持同步:

with pl.StringCache():
    bears = ...
    other_bears = ...
    print(
        bears.with_columns(
            pl.col("bears").to_physical().alias("physical")
        ),
        other_bears.with_columns(
            pl.col("bears").to_physical().alias("physical")
        )
    )

# Output
shape: (4, 2)
┌─────────┬──────────┐
│ bears   ┆ physical │
│ ---     ┆ ---      │
│ cat     ┆ u32      │
╞═════════╪══════════╡
│ Polar   ┆ 0        │
│ Grizzly ┆ 1        │
│ Sun     ┆ 2        │
│ Polar   ┆ 0        │
└─────────┴──────────┘ 
shape: (4, 2)
┌────────────┬──────────┐
│ bears      ┆ physical │
│ ---        ┆ ---      │
│ cat        ┆ u32      │
╞════════════╪══════════╡
│ Brown      ┆ 3        │
│ Spectacled ┆ 4        │
│ Polar      ┆ 0        │
│ Sloth      ┆ 5        │
└────────────┴──────────┘

可以看到,两个DataFrame中的Polar类别现在都用0表示,保持一致。StringCache方法让数值保持同步,从而在合并不同Series时不必重新编码Categorical。它解决了映射同步问题,但其内部实现仍有一些低效之处,而流式引擎使这些问题暴露出来。

不兼容之处

流式引擎按数据块(morsel)处理数据:将输入切分成小块,再在流水线中分别处理。这种方式的优势之一是,Polars可以在仍从磁盘或网络读取剩余数据时,就开始处理已经读到的数据。

在旧的Categorical局部实现中,映射按Series保存,所以每个数据块都有自己的局部映射,关联物理类型的值与字符串表示。Categorical的编码取决于输入顺序,而不同数据块中的值顺序几乎总是不同,因此局部映射也几乎总是不同步。

这些数据块重新组合为中间结果或最终结果时,必须不断重新编码,严重影响性能。

另一种选择——全局StringCache——同样承受着压力。由于Polars并行处理数据,许多数据块会同时执行。每次处理都要读取或更新全局StringCache中的映射,需要获取该对象的锁,导致流式流水线停顿。当多个Categorical共用一个全局StringCache时,影响更大,因为它们全部指向同一份映射。

这使Categorical操作变得非常缓慢。此外,StringCache也无法适用于分布式架构。

新的实现

为了与流式引擎兼容,我们吸取这些经验,推出了Categorical的新实现。

创建Categorical列时,现在可以传入pl.Categories()对象,定义物理类型与字符串表示之间的映射。

bear_categories =  pl.Categories(name="bears", namespace="org.polars", physical=pl.UInt8)
bears = (
    pl.DataFrame({
        "bears": ["Polar", "Grizzly", "Sun", "Polar"]
    })
    .cast(pl.Categorical(bear_categories))
)

Categories接受以下参数:

  • name:类别名称。
  • namespace:类别的命名空间。需要多个同名但类别不同的Categories时,可用它区分。
  • physical:类别的物理类型。可以使用默认的pl.UInt32(可编码超过42亿个类别)、pl.UInt16(原文列出的容量为65,535个类别),或pl.UInt8(原文列出的容量为255个类别)。选择合适的物理类型会影响应用的内存占用和性能。

它们的工作方式类似过去的全局StringCache,但有一些关键区别。

可以为它们设置名称与命名空间,让多份映射共存,也不再需要上下文包装器with pl.StringCache():。

只要名称、命名空间和底层物理类型相同,Categories就会匹配,即便是分别调用Categories创建的对象也一样。它们的值按字符串表示进行字典序(字母顺序)排序。没有显式提供Categories时,Categorical会使用全局Categories映射。这份全局映射与其他全局Categories对象共享类别,类似过去的StringCache。

底层会构建一个CategoricalMapping,它是关联物理值与字符串表示的双向映射,使用自定义字符串驻留机制,支持并行更新与读取。这意味着多个线程可以同时向映射加入新字符串,只要字符串不同,就不必相互等待。若字符串已存在于映射中,线程之间完全不会因此发生争用。

旧的StringCache实现无法做到上述两点。

此外,CategoricalMapping中的元素会自动被垃圾回收,避免没有数据引用时映射仍然只增不减。过去的StringCache若不手动释放,就会出现这种情况。

这次重新设计吸取了先前的经验,现在能够高效支持流式引擎。

Enum

Enum类型一直与Categorical密切相关。Categorical的类别是动态的:当Series出现尚未编码的新值时,类别会即时更新。

但Enum的类别预先定义,无法更改,因此具有性能优势。在新实现中,Enum使用FrozenCategoricalMapping,一次定义后便不可变。

此外,它允许定义类别的顺序。这样可以优化排序,因为排序使用物理表示,而不是开销更大的字符串表示。下面通过日志级别示例说明其工作方式:

log_levels = pl.Enum(["Critical", "Warning", "Info", "Debug"])
df = pl.DataFrame(
    {
        "log": ["Query took longer than usual", 
            "Finished downloading data", "Finished query", 
            "Result length: 50023", "Result had null values"],
        "level": ["Warning", "Info", "Info", "Debug", "Critical"]
    },
    schema={"log": pl.String, "level": log_levels}
)
df.sort("level")

# Output
shape: (5, 2)
┌──────────────────────────────┬──────────┐
│ log                          ┆ level    │
│ ---                          ┆ ---      │
│ str                          ┆ enum     │
╞══════════════════════════════╪══════════╡
│ Result had null values       ┆ Critical │
│ Query took longer than usual ┆ Warning  │
│ Finished downloading data    ┆ Info     │
│ Finished query               ┆ Info     │
│ Result length: 50023         ┆ Debug    │
└──────────────────────────────┴──────────┘

总结

新的Categorical实现解决了旧设计与Polars流式引擎及分布式架构不兼容的核心性能瓶颈。通过引入具有无争用字符串驻留和自动垃圾回收能力的Categories对象,我们消除了旧实现中昂贵的重新编码操作与锁争用。

在分布式查询中,重新编码仍然不可避免,但Categorical依然可能值得使用,因为它能大幅减少shuffle的数据量。不过,如果事先知道全部类别,仍然建议优先使用Enum。

这次重新设计让你更好地控制流水线中Categorical的管理方式。现在可以定义具有特定命名空间与物理类型的命名类别,更容易分析内存占用和性能。如果处理的是预先定义、具有顺序的类别,Enum提供了性能更高的选择,我们建议使用它。

更新后,StringCache已经成为不执行实际操作的no-op;由于Polars默认使用全局Categories对象,原文指出不会产生行为变化。如果代码当前使用StringCache,迁移方式是:在需要的位置,用显式的Categories对象替换上下文管理器。

这些改进让你可以在流式操作中放心使用Categorical与Enum,无需担心旧实现带来的性能退化,使它们成为大规模数据处理流水线中的实用选择。

原文:Understanding the New Categorical,作者 Thijs Nieuwdorp,2026年1月29日,Polars。本文为中文翻译。示例输出及性能判断来自原文;原文以 Polars 1.31.0 演示旧行为,后续示例使用重构后的 API。含省略号的片段是示意代码,使用时需补全。文中的第一人称指原作者及 Polars 团队。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容