Spark DataFrame过滤异常:生产环境无记录本地却有结果
问题原因分析
核心问题出在Spark的惰性计算机制和自定义UDF中随机数生成的非确定性:
- Spark的DataFrame是惰性求值的,
sampleBasedAll不会立即执行计算,只有当触发write(第一个Action操作)和count()(第二个Action操作)时,才会重新运行整个DAG逻辑。这意味着两次Action会分别执行一遍UDF,生成的sampled列值可能完全不同。 - 你UDF中使用的
random实例如果在Driver端初始化,在分布式环境下会被序列化到各个Executor的Task中。随机数生成器的状态在序列化/反序列化过程中可能出现异常,或者不同Task、不同计算阶段的随机数序列完全不一致。生产环境中,第一次write时生成了包含sampled=1的记录,但第二次count()时重新执行UDF,恰好生成的sampled全为0,所以计数为0;而本地是读取已经持久化的CSV文件,sampled值固定,自然能查到结果。
解决方法
方法一:改用Spark内置采样函数(推荐)
Spark提供了分布式安全的分层采样API,替代自定义UDF避免非确定性问题:
// 按count列定义分层采样比例 val fractions = df.select("count").distinct().map { row => val count = row.getInt(0) val fraction = if (count < TEN_K_SELLER_ITEM_BENCH) 1.0 else 10000.0 / count (count, fraction) }.collectAsMap() // 执行分层采样,固定种子保证结果可复现 val sampledDF = df.sampleBy("count", fractions, seed = 1234L) scribe.info("sampledDF.count = " + sampledDF.count())
方法二:持久化DataFrame
在生成sampleBasedAll后调用cache(),让Spark缓存计算结果,两次Action复用同一批数据:
val sampleBasedAll = df.withColumn("sampled", sampledOrNot(col("count"))).cache() // 先触发缓存计算 sampleBasedAll.count() // 后续write和count均使用缓存数据 sampleBasedAll.repartition(10).write.option("header", value = true).option("compression", "gzip").csv("/sampleBasedAll") val sampledDF = sampleBasedAll.repartition(100).filter("sampled = 1").select($"sellerId", $"siteId", $"count", $"desc") scribe.info("sampledDF.count = " + sampledDF.count())
方法三:修复UDF中的随机数生成
若必须使用自定义UDF,需在UDF内部初始化随机数生成器并绑定固定种子,保证分布式环境下的一致性:
def sampledOrNot = udf((count: Int) => { if(count < TEN_K_SELLER_ITEM_BENCH){ 1 }else{ // 基于线程ID初始化随机数生成器,保证每个Task的随机性 val random = new scala.util.Random(Thread.currentThread().getId) val targetValue = 10000/count.toDouble var base = 1 var adjustedTarget = targetValue while (adjustedTarget < 1){ adjustedTarget *= 10 base *= 10 } val randomId = random.nextLong(0, 1000000000000L) if(randomId % base <= (adjustedTarget.intValue() + 1)) 1 else 0 } })
内容的提问来源于stack exchange,提问作者LizzyMM
相关产品推荐
相关产品推荐

