Spark中如何可靠打印分区长度?遇Ordinal must >=1报错求助
如何在Spark 2.3.0(Python)中可靠获取每个分区的记录数
嘿,我来帮你解决这个问题!你遇到的Ordinal must >= 1错误大概率是代码变体里误操作了分区索引的序号规则,先给你两个能稳定运行的方案,再聊聊错误原因。
方案1:修正mapPartitionsWithIndex的写法
你的核心思路没问题,试试这个经过验证的版本,它能准确返回每个分区的索引(从0开始)和对应行数:
# 获取每个分区的索引和记录数 partition_stats = df.rdd.mapPartitionsWithIndex( lambda partition_idx, iter: [(partition_idx, sum(1 for _ in iter))] ).collect() # 逐个打印结果 for idx, count in partition_stats: print(f"分区 {idx} 的行数: {count}")
这个代码在Spark 2.3.0环境下可以直接运行,不会触发序号相关的错误。
方案2:用glom()实现更简洁的写法
如果只是需要分区大小,glom()方法会更简洁——它把每个分区的元素转为列表,直接取长度就行:
# 获取所有分区的长度,再加上索引 partition_sizes = df.rdd.glom().map(len).collect() for idx, size in enumerate(partition_sizes): print(f"分区 {idx} 的行数: {size}")
这个写法代码更短,同样适合Spark 2.3.0,且不容易踩索引的坑。
关于Ordinal must >= 1错误的分析
这个错误通常是因为你在某个代码变体里,尝试将分区索引强制转为从1开始的序号(比如写了partition_idx + 1),但后续逻辑(或者Spark内部的某些校验)要求序号必须≥1,而如果存在空分区或者索引处理异常,就会触发这个报错。Spark的分区索引默认是从0开始的,除非业务明确要求,否则不要随意修改这个规则。
内容的提问来源于stack exchange,提问作者Sean Lindo
相关产品推荐
相关产品推荐

