如何在Scala Spark的substring函数中传入同行列值?
问题解决方法
Spark中直接调用Scala的substring函数时,第三个参数要求传入字面量Int,无法直接传入Column类型的列值。要实现用另一列的值作为字符串截取长度,有两种实用方案:
方案1:使用Column自带的substring方法
Spark的Column对象内置了支持列参数的substring方法,写法如下:
val readyDF: Dataset[Row] = peopleWithJobsAndAgeLimitsDF.withColumn("fnamelname", col("fnamelname").substring(0, col("ageLimit")) )
注:从你的Schema来看ageLimit本身就是Integer类型,无需额外cast("Int")转换。
方案2:用expr写SQL风格表达式
如果更习惯SQL语法,可以通过expr函数直接编写SQL逻辑,天然支持引用列作为参数:
import org.apache.spark.sql.functions.expr val readyDF: Dataset[Row] = peopleWithJobsAndAgeLimitsDF.withColumn("fnamelname", expr("substring(fnamelname, 0, ageLimit)") )
可选:处理ageLimit为0的场景
当ageLimit为0时,上述代码会返回空字符串,若需要保留原字段内容,可以用when函数增加判断逻辑:
import org.apache.spark.sql.functions.when val readyDF: Dataset[Row] = peopleWithJobsAndAgeLimitsDF.withColumn("fnamelname", when(col("ageLimit") > 0, col("fnamelname").substring(0, col("ageLimit"))) .otherwise(col("fnamelname")) )
内容的提问来源于stack exchange,提问作者rafald
相关产品推荐
相关产品推荐

