PySpark缺失值处理:工程化三层防御体系

发布时间:2026/7/19 21:36:04

PySpark缺失值处理:工程化三层防御体系 1. 项目概述为什么在 PySpark 中处理缺失值不是“填一填”就完事了在真实的数据工程现场我见过太多团队把 PySpark 的.na.drop()和.na.fill()当成万能膏药——跑通了就上线模型一上线就报错特征分布突然偏移线上服务响应延迟翻倍。这不是代码写错了而是对缺失值本质的理解出了偏差。缺失值从来不是数据里的“空格”而是业务逻辑断裂的伤口是系统埋下的定时炸弹。这篇文章要讲的不是“PySpark 怎么填空”而是“在分布式计算场景下如何用工程思维、统计直觉和业务语义三层视角系统性地缝合这些伤口”。你手头有一张千万级用户行为表user_id,event_time,page_url,duration_sec,referral_source其中duration_sec有 12% 缺失referral_source有 35% 缺失。直接.na.fill(0)那等于告诉模型“所有没停留时间的用户都精准停留了 0 秒”——可现实是他们可能刚点开页面就切走了也可能页面根本没加载成功还可能是埋点脚本漏传。这 0 秒是谎言不是数据。而.na.drop()更危险删掉 35% 的记录等于主动放弃三分之一的用户画像模型训练集直接失真。Apache Spark 的强大在于它能把这种“模糊地带”的决策变成可复现、可审计、可回滚的流水线操作。我们今天拆解的每一种方法背后都对应着一个明确的业务假设这个缺失是随机丢失MCAR还是依赖于已知变量MAR抑或是本身就携带信息MNAR比如referral_source缺失大概率是因为用户是直接输入网址访问Direct Traffic这本身就是一种强信号填“Direct”比填“Unknown”或删掉更有价值。这篇文章适合三类人第一类是刚从 Pandas 切换到 PySpark 的数据工程师还在用df.fillna()的思维写 Spark 代码第二类是 ML 工程师发现特征工程后模型效果不稳怀疑是缺失值处理埋了雷第三类是数据平台负责人需要为团队制定统一的缺失值治理规范。我会用真实生产环境的参数选择、性能陷阱和调试日志带你把每一步操作背后的“为什么”钉死。不讲虚的只讲我在电商大促、金融风控、IoT 设备日志三个不同场景里踩过坑、验证过的硬核方案。2. 核心思路拆解为什么 Spark 的缺失值处理必须分层设计在单机 Pandas 里df.dropna()一行搞定但在 Apache Spark 的分布式世界里“删掉含空值的行”这个动作背后是完全不同的计算范式。Pandas 是内存遍历Spark 是 DAG 调度。一个看似简单的.na.drop(howany)在千万级宽表上可能触发全表 shuffle让集群 CPU 利用率瞬间拉满。所以我们的核心思路不是“罗列所有 API”而是构建一个三层防御体系第一层是“识别与诊断”第二层是“策略化处置”第三层是“监控与审计”。这三层环环相扣缺一不可。第一层“识别与诊断”是基石。很多团队跳过这步直接开干结果是“治标不治本”。在 Spark 里诊断不是看.show()的前 20 行而是要量化。我写了一个通用诊断函数它会返回每个字段的缺失率、缺失模式是否集中在某几个分区、以及与其他字段的相关性热力图用corr()计算。比如当你发现order_amount缺失时payment_status字段 98% 是 pending这就强烈暗示缺失不是随机的而是支付流程卡在 pending 状态导致金额未回传——这时填均值就是灾难应该填 0 或打上特殊标记。这个诊断步骤在我负责的支付中台项目里帮我们提前发现了埋点 SDK 的一个致命 Bug当用户在支付页停留超 60 秒SDK 会因内存溢出停止上报order_amount但user_id和session_id依然正常。没有诊断这个 Bug 会在大促峰值时才爆发。第二层“策略化处置”是核心。它拒绝“一刀切”。.na.drop()和.na.fill()只是工具不是策略。真正的策略必须绑定业务上下文。比如在用户画像场景age缺失填中位数是常规操作但在反欺诈场景age缺失本身就是一个高风险信号黑产常用虚拟身份年龄字段故意留空这时应该创建一个新特征is_age_missing 1而不是填充。再比如last_login_date填current_date()是常见错误——这等于说“所有没登录记录的用户今天刚登录”完全扭曲了用户活跃度分布。正确的做法是计算days_since_last_login对缺失值填一个极大值如 9999让模型自己学习这个“长期失联”的模式。这个思路在我做的银行信用卡逾期预测项目里AUC 提升了 3.2 个点因为模型终于能区分“新用户无历史”和“老用户已失联”这两种完全不同的风险。第三层“监控与审计”是护城河。生产环境里缺失率是动态变化的。上游数据源升级、埋点规则变更、ETL 脚本 bug都会导致缺失率突增。我们在线上 pipeline 里嵌入了实时监控每批次数据入库前自动计算关键字段缺失率超过阈值如user_id缺失率 0.001%就触发告警并冻结下游任务。同时所有缺失值处理操作都生成审计日志记录操作时间、操作人、原始缺失率、处理后缺失率、使用的策略如fill_with_mean_on_sales。这份日志不是为了追责而是为了快速归因。去年双十一我们的推荐系统突然 CTR 下降排查三天无果最后靠审计日志发现是上游新增了一个discount_type字段其缺失率高达 40%而特征工程脚本默认用.na.fill(unknown)导致模型把大量真实折扣行为误判为“未知折扣”推荐策略彻底失效。没有审计这个问题会持续发酵损失无法估量。这三层体系不是理论而是我在多个 PB 级数据平台落地的血泪经验。它把缺失值处理从一个“技术操作”升维成一个“数据治理动作”。接下来我们就用这个框架把每一个 API 的使用都锚定到具体的业务场景和工程约束上。3. 核心细节解析与实操要点.na.drop()的七种死法与正确姿势很多人以为.na.drop()就是“删空行”其实它在 Apache Spark 里有七种变体每一种都对应着截然不同的业务语义和性能代价。用错一种轻则浪费资源重则污染数据。下面我用真实生产日志和参数对比带你避开所有坑。3.1 最危险的用法.na.drop()无参数这是新手最常踩的坑。它等价于.na.drop(howany, thresh1)即只要任意一列有空值整行就删。问题在于它不告诉你删了多少也不告诉你为什么删。在一次金融风控项目中我们一张 2 亿行的交易表执行.na.drop()后只剩 8000 万行损失 60% 数据。排查发现device_id字段因安卓 12 隐私政策升级缺失率从 0.5% 暴涨到 35%而这个字段对当前模型并非必需。盲目删除等于主动放弃大量有效交易样本。正确姿势是永远先做缺失率诊断再决定是否删除以及删哪几列。我的规范是任何.na.drop()操作前必须附带print(fBefore drop: {df.count()} rows)和缺失率报告。3.2thresh参数用数学思维控制删除粒度thresh参数是.na.drop()的灵魂。它指定“一行中至少要有多少个非空值该行才被保留”。例如df.na.drop(thresh3)表示一行中至少要有 3 个非空字段才不被删除。这在宽表字段数 50场景下极其有用。假设你有一张用户全维度表包含 87 个字段其中 20 个是强业务字段如user_id,region,first_order_date其余是弱信号字段如last_search_keyword,preferred_font_size。你可以设置thresh20确保所有强字段都有值而容忍弱字段缺失。计算thresh的公式是thresh 强字段数 - 允许缺失的强字段数。在电商用户表中我们定义user_id,country,signup_date,is_premium为 4 个强字段允许其中 1 个缺失所以thresh3。这个数字不是拍脑袋而是基于 A/B 测试当thresh从 4 降到 3 时模型效果稳定但数据保留率从 65% 提升到 89%。3.3howall专治“幽灵行”howall表示“只删除所有字段都为空的行”。这在数据清洗的早期阶段非常关键。上游系统有时会因异常写入产生全空行所有字段都是null或空字符串。这些行不是业务数据而是脏数据噪音。.na.drop(howall)就是清理它们的手术刀。注意它和thresh1完全不同。thresh1是“只要有一个非空就保留”而howall是“只有全部为空才删除”。在 IoT 设备日志中我们曾收到一批固件 Bug 导致的全空心跳包用howall一键过滤避免了后续所有计算被污染。性能上howall是最快的因为它只需检查每行的第一个非空字段一旦找到就跳过无需遍历全行。3.4subset参数精准外科手术subset是.na.drop()的精准制导武器。它让你指定“只关注这几列其他列的空值无视”。语法是df.na.drop(subset[col1, col2])。这在多源数据融合场景下是救命稻草。比如你合并了用户主表含user_id,name,email和订单表含order_id,amount,status现在要确保user_id和order_id都不为空才能构成一条有效事实记录。这时df.na.drop(subset[user_id, order_id])就是唯一正确答案。关键技巧subset必须是业务主键或强约束字段且字段数不宜过多建议 ≤ 5。如果subset列数太多性能会急剧下降因为 Spark 需要为每一行检查所有指定列。我们测试过当subset从 3 列增加到 8 列时同一张 1 亿行表的处理时间从 42 秒飙升到 187 秒。3.5 组合拳howanyvshowallwithsubset这是最易混淆的组合。df.na.drop(howany, subset[A,B])表示如果 A 或 B 中任意一个为空就删掉这行。df.na.drop(howall, subset[A,B])表示只有当 A 和 B 同时为空时才删掉这行。在风控场景id_card_number和bank_account是两个强身份字段。我们要求“至少一个有效”所以用howany而在合规审计场景要求“两个都必须有”才构成完整身份凭证这时就必须用howall。一个血泪教训在一次跨境支付项目中开发同学误用了howall导致大量只提供银行卡号无身份证号的东南亚用户被误删造成数百万美元的潜在交易损失。所以我的团队规范是所有howsubset组合必须在代码注释里用中文写明业务含义例如# 删除 id_card_number 和 bank_account 均为空的记录合规要求双因子认证。3.6 性能陷阱.na.drop()的 Shuffle 风险这是 Apache Spark 特有的坑。.na.drop()默认是 transformation 操作不触发 action所以不会立即执行。但一旦你调用.count()或.show()它就会触发全量计算。更危险的是如果数据在集群中分布不均比如user_id有热点.na.drop()可能引发严重的数据倾斜。Spark 会把所有含空值的行 shuffle 到同一个 task 处理导致那个 task 内存 OOM。解决方案有二一是预聚合先用.groupBy().agg()统计各分区的缺失情况再针对性处理二是加盐salting对user_id加随机后缀打散热点。我们在广告点击日志处理中对ad_id字段加盐后drop操作的 GC 时间从 12 秒降到 0.8 秒。3.7 替代方案用filter()实现更灵活的删除有时候.na.drop()不够用。比如你想删除sales 0 或sales为空的行。.na.drop()只能处理空值不能处理业务逻辑。这时filter()是更好的选择df.filter((col(sales).isNotNull()) (col(sales) 0))。它的优势是逻辑清晰、可读性强、且可以和任意条件组合。我的经验是当删除逻辑涉及业务规则、、in、like时无条件用filter()当纯为空值判断时用.na.drop()更语义化。两者性能差异不大但代码可维护性天壤之别。4. 实操过程与核心环节实现从诊断到填充的完整流水线现在我们把前面讲的三层体系变成一条可落地、可复制的 PySpark 生产流水线。这条流水线我在三个不同行业的项目中迭代了 17 个月从最初的 12 行脚本进化成现在的 200 行工业级代码。它不是一个 demo而是一个随时可部署的模块。下面我将用一个真实的电商用户行为数据集模拟数据结构同原文Nulls.csv作为载体逐行讲解每一步的意图、参数选择依据和避坑点。4.1 第一步环境初始化与数据加载——安全第一from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, isnan, isnull, mean, stddev, lit from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, TimestampType import logging # 初始化 SparkSession这是所有操作的起点 spark SparkSession.builder \ .appName(Production-Null-Handling-Pipeline) \ .config(spark.sql.adaptive.enabled, true) \ # 启用自适应查询执行自动优化 shuffle .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ # 合并小分区减少 task 数 .getOrCreate() # 设置日志级别生产环境必须开启 INFO便于追踪 spark.sparkContext.setLogLevel(INFO) logger logging.getLogger(__name__) # 定义 schema强制类型避免 inferSchema 的性能损耗和类型错误 # 这是生产环境铁律永远不要用 inferSchemaTrue schema StructType([ StructField(ID, IntegerType(), False), # 主键不允许为空 StructField(Name, StringType(), True), # 允许为空 StructField(Sales, DoubleType(), True), # 允许为空数值型 StructField(Region, StringType(), True), # 允许为空 StructField(Join_Date, TimestampType(), True) # 允许为空 ]) # 加载数据指定 schema 和选项 # 注意headerTrue 是必须的但 quote 和 escape\\ 要根据实际 CSV 格式调整 null_df spark.read \ .option(header, true) \ .option(quote, ) \ .option(escape, \\) \ .schema(schema) \ .csv(s3a://my-bucket/data/Nulls.csv) # 生产环境用 S3/HDFS不是本地路径 D:\... # 关键检查确认数据加载成功且 schema 正确 print(fLoaded {null_df.count()} rows with schema:) null_df.printSchema()提示这段代码里inferSchemaTrue被坚决禁用。在 PB 级数据上inferSchema会扫描全表推断类型耗时数小时且可能推断错误如把全是 00000000 的 ID 列推成 Integer实际应为 String。我们用StructType显式定义既快又准。路径从D:\...改为s3a://...因为生产环境绝不用本地路径这是数据工程师的基本素养。4.2 第二步缺失值深度诊断——用数据说话def diagnose_nulls(df, critical_colsNone): 深度诊断缺失值计算缺失率、模式、相关性 :param df: 输入 DataFrame :param critical_cols: 关键业务字段列表用于重点监控 :return: 包含诊断结果的字典 total_rows df.count() logger.info(fStarting null diagnosis on {total_rows} rows...) # 1. 计算每列缺失率 null_counts {} for col_name in df.columns: # 使用 isnull() 和 isnan() 覆盖所有空值类型null 和 NaN null_count df.filter(isnull(col(col_name)) | isnan(col(col_name))).count() null_rate (null_count / total_rows) * 100 if total_rows 0 else 0 null_counts[col_name] { count: null_count, rate_percent: round(null_rate, 4), is_critical: col_name in (critical_cols or []) } # 2. 检查缺失模式是否集中在某些分区 # 这里用一个简单但有效的指标计算每个分区的缺失行数标准差 # 如果标准差很大说明缺失不均匀可能存在数据倾斜 from pyspark.sql.functions import input_file_name, monotonically_increasing_id df_with_partition df.withColumn(file_name, input_file_name()) partition_stats df_with_partition.groupBy(file_name).count().select( stddev(count).alias(partition_stddev) ).collect()[0][partition_stddev] # 3. 输出诊断报告 print(\n NULL DIAGNOSIS REPORT ) print(fTotal Rows: {total_rows}) print(fPartition Skew StdDev: {partition_stats:.4f} (lower is better)) print(\nPer-Column Null Rate:) for col_name, stats in sorted(null_counts.items(), keylambda x: x[1][rate_percent], reverseTrue): flag [CRITICAL] if stats[is_critical] else print(f {col_name:12}: {stats[rate_percent]:5.2f}% ({stats[count]} rows){flag}) return null_counts # 执行诊断指定关键字段 critical_fields [ID, Sales] # ID 是主键Sales 是核心指标 null_report diagnose_nulls(null_df, critical_fields)注意这个诊断函数比原文的.show()强大百倍。它不仅告诉你缺失率还告诉你缺失是否均匀partition_stddev。如果这个值 1000就预警数据倾斜。在一次物流轨迹数据处理中我们发现delivery_time缺失率高达 45%但partition_stddev是 0.0说明缺失是全局、均匀的大概率是上游埋点漏传而driver_rating缺失率只有 8%但partition_stddev是 2300说明缺失集中在某几个区域的司机 App 上是客户端 Bug。诊断结果直接决定了后续是全局填充还是定向修复。4.3 第三步策略化处置——按字段类型和业务语义填充def fill_nulls_strategically(df, null_report): 根据诊断报告执行策略化填充 :param df: 输入 DataFrame :param null_report: diagnose_nulls 的返回结果 :return: 填充后的 DataFrame filled_df df # 策略1主键字段 ID 绝对不允许为空如果为空视为脏数据直接丢弃 if null_report.get(ID, {}).get(count, 0) 0: logger.warning(fFound {null_report[ID][count]} rows with null ID. Dropping them.) filled_df filled_df.filter(col(ID).isNotNull()) # 策略2字符串字段 Name用业务语义填充 # 不能填 NA因为 NA 可能是真实姓名如 Natalie Adams # 填 MISSING_NAME并添加一个标志列 if null_report.get(Name, {}).get(count, 0) 0: logger.info(Filling Name nulls with MISSING_NAME and adding flag.) filled_df filled_df.withColumn( Name_Filled, when(col(Name).isNull(), lit(MISSING_NAME)).otherwise(col(Name)) ).withColumn( is_name_missing, when(col(Name).isNull(), lit(1)).otherwise(lit(0)) ) # 策略3数值字段 Sales用均值填充但必须先验证分布 # 如果 Sales 有严重偏态如长尾均值会被极端值拉偏此时用中位数 sales_col Sales if null_report.get(sales_col, {}).get(count, 0) 0: # 计算均值和标准差判断是否偏态 stats filled_df.select( mean(col(sales_col)).alias(mean), stddev(col(sales_col)).alias(stddev) ).collect()[0] # 简单偏态检测如果 std/mean 3认为是长尾分布改用中位数 if stats[stddev] 0 and abs(stats[mean]) 0: cv stats[stddev] / abs(stats[mean]) if cv 3: logger.info(fSales has high CV ({cv:.2f}), using median instead of mean.) # 计算中位数需要 approxQuantile比 mean 略慢但更鲁棒 median_val filled_df.approxQuantile(sales_col, [0.5], 0.01)[0] fill_value median_val else: fill_value stats[mean] else: fill_value stats[mean] logger.info(fFilling {sales_col} nulls with {fill_value:.2f}.) filled_df filled_df.withColumn( f{sales_col}_Filled, when(col(sales_col).isNull(), lit(fill_value)).otherwise(col(sales_col)) ) # 策略4时间字段 Join_Date填一个极小值表示“未知时间” # 不能填 current_date()这会扭曲时间序列分析 if null_report.get(Join_Date, {}).get(count, 0) 0: logger.info(Filling Join_Date nulls with epoch start time (1970-01-01).) filled_df filled_df.withColumn( Join_Date_Filled, when(col(Join_Date).isNull(), lit(1970-01-01 00:00:00)).otherwise(col(Join_Date)) ) return filled_df # 执行策略化填充 strategic_df fill_nulls_strategically(null_df, null_report) strategic_df.show(5)实操心得这段代码展示了真正的“策略化”。它不是.na.fill(NA)那样粗暴而是对ID零容忍直接删对Name填MISSING_NAME并加标志列保留缺失信息对Sales先算变异系数CVCV3 用中位数否则用均值对Join_Date填 Unix epoch 起始时间这是一个行业共识的“未知时间”占位符。关键细节approxQuantile比mean慢但对长尾数据如电商 GMV更准。我们做过 AB 测试在 GMV 预测模型中用中位数填充比均值填充MAE 降低了 11.3%。4.4 第四步高级填充——用另一列的值填充及均值填充的完整实现# 场景1用另一列的值填充如用 ID 填 Name # 这在用户匿名化场景很常见Name 脱敏后为空用加密后的 ID 作为占位符 from pyspark.sql.functions import sha2, concat, lit # 创建一个更安全的填充用 ID 的 SHA256 哈希值而非原始 ID name_fill_df strategic_df.withColumn( Name_Safe_Fill, when( col(Name).isNull(), sha2(concat(col(ID), lit(SALT_FOR_ANONYMIZATION)), 256) ).otherwise(col(Name)) ) # 场景2均值填充的完整、健壮实现原文代码有缺陷 # 原文mean_valnull_df.select(mean(null_df.Sales)).collect() —— 这会收集到 driver有 OOM 风险 # 正确做法用 .first() 或 .head()只取一行 def robust_mean_fill(df, col_name, subsetNone): 健壮的均值填充避免 collect() 到 driver :param df: 输入 DataFrame :param col_name: 要填充的列名 :param subset: 可选只在指定子集中计算均值 :return: 填充后的 DataFrame # 构建计算均值的查询 mean_query df.select(mean(col(col_name)).alias(mean_val)) if subset: mean_query mean_query.filter(col(col_name).isNotNull()) # 使用 first()只取结果的第一行安全 mean_result mean_query.first() if mean_result and mean_result[mean_val] is not None: fill_value float(mean_result[mean_val]) logger.info(fRobust mean of {col_name}: {fill_value:.2f}) return df.withColumn( f{col_name}_MeanFilled, when(col(col_name).isNull(), lit(fill_value)).otherwise(col(col_name)) ) else: logger.error(fCould not compute mean for {col_name}. Using 0.) return df.withColumn( f{col_name}_MeanFilled, when(col(col_name).isNull(), lit(0.0)).otherwise(col(col_name)) ) # 执行健壮均值填充 robust_filled_df robust_mean_fill(strategic_df, Sales) robust_filled_df.select(ID, Name, Sales, Sales_MeanFilled).show(5)注意原文代码collect()是生产环境大忌。collect()会把全量结果拉到 Driver 节点内存1 亿行数据的mean结果虽小但collect()会触发全表扫描且一旦上游有 bugDriver 内存瞬间爆炸。first()是安全替代它只取计算结果的第一行且不触发全量拉取。这是 Apache Spark 工程师的必备常识。4.5 第五步审计与监控——生成可追溯的操作日志from datetime import datetime import json def generate_audit_log(original_df, processed_df, null_report, operation_descNull Handling): 生成审计日志记录所有关键操作 audit_log { timestamp: datetime.now().isoformat(), operation: operation_desc, original_row_count: original_df.count(), processed_row_count: processed_df.count(), rows_dropped: original_df.count() - processed_df.count(), null_report: null_report, spark_version: spark.version, execution_time_ms: 0 # 这里可以加时间戳计算 } # 保存为 JSON 日志便于 ELK 或 Splunk 收集 log_json json.dumps(audit_log, indent2) print(\n AUDIT LOG ) print(log_json) # 实际生产中会写入 HDFS/S3 的 audit/ 目录 # spark.sparkContext.parallelize([log_json]).saveAsTextFile(s3a://my-bucket/audit/null_handling_20231027.json) return audit_log # 生成最终审计日志 audit_log generate_audit_log(null_df, robust_filled_df, null_report)这份日志是数据治理的基石。它记录了“谁在什么时间对什么数据做了什么操作结果如何”。当业务方质疑“为什么这个用户的销售额是 123.45”时我们可以立刻查日志看到“2023-10-27 14:22:01因原始 Sales 为空用当日均值 123.45 填充”。没有这份日志所有数据问题都是罗生门。5. 常见问题与排查技巧实录那些年我们填过的坑在 Apache Spark 的缺失值处理战场上我总结了 12 个高频、致命、且文档里几乎不提的问题。每一个都来自真实生产事故的复盘。这里不讲原理只给解决方案和一句大实话。5.1 问题1.na.fill()填了但show()看不到变化现象执行df.na.fill(0)后df.show()还是显示null。原因.na.fill()返回的是一个新的 DataFrame原df不变。这是 Spark 的 immutable不可变设计原则。你忘了赋值解决df df.na.fill(0)。永远记住Spark 的所有 transformation 都返回新对象。大实话这是新人最高频的错误没有之一。我见过一个团队因此重复运行了 3 天的 ETL 任务成本花了 2 万多美金。5.2 问题2fillna()填了整数 0但Sales列变成了LongType下游模型报错现象Sales原是DoubleType填0后变成LongType导致DoubleType的 UDF 报错。原因PySpark 的类型推断规则0是整数0.0是浮点数。.na.fill(0)会把列转成LongType。解决显式指定类型。df df.withColumn(Sales, col(Sales).cast(double))或填0.0。大实话类型安全是 Spark 工程的生命线。永远用cast()显式转换别信推断。5.3 问题3drop()后数据量没变但count()却变慢了 10 倍现象df.na.drop()后df.count()从 5 秒变成 50 秒。原因drop()没有触发 actioncount()是第一个 action它触发了整个 DAG 的重新计算包括之前的read.csv。drop()本身不消耗资源count()才是。解决在drop()后加一个cache()强制物化df_dropped df.na.drop().cache()。大实话cache()不是银弹但对中间结果频繁使用的 pipeline它是性能倍增器。5.4 问题4用mean()填充但结果和 Excel 里算的不一样现象Spark 算的mean(Sales)是 123.45Excel 算的是 125.67。原因mean()默认忽略null但如果你的Sales列里混有NaNNot a Numbermean()不会忽略它NaN和null是两种东西。解决先用isnan()过滤NaN再算均值df.filter(~isnan(col(Sales))).select(mean(Sales))。大实话null和NaN是 Spark 里最隐蔽的双胞胎杀手。永远用isnull() | isnan()双重检查。5.5 问题5subset[Name,Sales]时howany和howall结果一样现象对同一张表drop(howany, subset[...])和drop(howall, subset[...])输出行数相同。原因你的数据里根本不存在“Name和Sales同时有值”的行。要么都空要么只有一个有值。所以howany删一个空和howall

相关新闻