Scala Spark:访问DataFrame的Map键值、生成score表及存储score到变量
嘿,我来帮你搞定这两个Spark Scala的操作!先理清楚你的表结构:id是字符串类型,persona_model是嵌套的map<string, struct<score:double, tag:string>>,还有分区字段process_date。下面分两步逐个实现:
1. 生成仅包含id和对应score的新表
首先读取原表,然后根据你的需求提取id和对应的score字段。这里分两种常见情况:
情况1:Map中只有单个固定Key(比如示例中的"Tech")
如果你的persona_model里始终只有一个固定的Key(比如示例里的"Tech"),可以直接通过键索引提取score:
// 读取原表,假设表名为persona_table val originalDF = spark.table("persona_table") // 选择id和对应score,并给score字段起别名 val idScoreDF = originalDF.select( col("id"), col("persona_model")("Tech")("score").alias("score") ) // 将结果保存为新表(根据需求选择存储模式,这里用overwrite覆盖已有表) idScoreDF.write.mode("overwrite").saveAsTable("id_score_table")
情况2:Map中有多个Key,需要展开所有键值对
如果persona_model里有多个分类Key,需要把每个Key对应的score都展开,可以用explode函数:
import org.apache.spark.sql.functions.explode val originalDF = spark.table("persona_table") // 先展开Map的键值对,再提取score val explodedIdScoreDF = originalDF.select( col("id"), explode(col("persona_model")).alias("category", "persona_struct") ).select( col("id"), col("persona_struct.score").alias("score"), col("category") // 可选:保留分类名称 ) // 保存为新表 explodedIdScoreDF.write.mode("overwrite").saveAsTable("id_score_all_categories_table")
2. 访问特定score值并存储到整数变量temp中
要把特定行的score转成整数并存到变量里,需要先过滤出目标行,再提取值。这里推荐两种方式,兼顾效率和安全性:
方式1:直接提取(假设目标行一定存在)
比如要获取示例中id为000000E91010441BB122402A45D439E7的score,并转成整数:
val temp: Int = originalDF .filter(col("id") === "000000E91010441BB122402A45D439E7") .select(col("persona_model")("Tech")("score").cast("int")) .first() // 获取匹配的第一行 .getInt(0) // 提取整数值
方式2:安全处理空值(防止没有匹配行报错)
如果不确定是否存在目标行,用headOption避免空指针异常,还可以设置默认值:
// 先获取可选的结果 val tempOption: Option[Int] = originalDF .filter(col("id") === "000000E91010441BB122402A45D439E7") .select(col("persona_model")("Tech")("score").cast("int")) .headOption() .map(_.getInt(0)) // 给变量赋值,没有匹配行时用默认值0 val temp: Int = tempOption.getOrElse(0)
内容的提问来源于stack exchange,提问作者Romal Jaiswal
相关产品推荐
相关产品推荐

