如何在Spark SQL中读取键值对?并用Spark SQL/Scala拆分表格列
问题1:如何在Spark SQL中读取键值对
假设你要读取的是文本格式的键值对数据(每行格式如key=value),可按以下步骤操作:
- 读取文本并转换为DataFrame
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().appName("ReadKVData").getOrCreate() // 读取文本文件,拆分每行的键值对 val kvRDD = spark.sparkContext.textFile("/path/to/kv-file.txt") .map(line => { val parts = line.split("=", 2) // 限制拆分次数,避免value含=符号时出错 (parts(0), parts(1)) }) // 转成DataFrame并注册临时视图 import spark.implicits._ val kvDF = kvRDD.toDF("key", "value") kvDF.createOrReplaceTempView("key_value_table")
- 用Spark SQL查询数据
SELECT key, value FROM key_value_table WHERE key = "target_key"
如果是读取HBase这类键值存储,可借助Spark的HBase连接器,配置表名、列族等参数后读取为DataFrame,再通过Spark SQL操作。
问题2:拆分表格列为多列
假设目标列的格式为"name:Alice,age:30,city:NewYork"这类键值对拼接字符串,以下是两种实现方式:
方式1:Spark SQL内置函数实现
先将源表注册为临时视图(假设表名为source_table,目标列名为info):
-- 固定顺序拆分 SELECT split(split(info, ",")[0], ":")[1] AS name, cast(split(split(info, ",")[1], ":")[1] AS int) AS age, split(split(info, ",")[2], ":")[1] AS city FROM source_table
如果键的顺序不固定,推荐使用str_to_map函数(Spark 2.3+支持):
SELECT kv_map['name'] AS name, cast(kv_map['age'] AS int) AS age, kv_map['city'] AS city FROM ( SELECT str_to_map(info, ",", ":") AS kv_map FROM source_table ) t
方式2:Scala代码实现
import org.apache.spark.sql.functions._ val sourceDF = spark.read.table("source_table") // 转换为Map类型后提取字段 val resultDF = sourceDF .withColumn("kv_map", str_to_map(col("info"), ",", ":")) .select( col("kv_map")("name").alias("name"), col("kv_map")("age").cast("int").alias("age"), col("kv_map")("city").alias("city") ) // 查看结果或写入表 resultDF.show() resultDF.write.saveAsTable("split_result_table")
如果目标列是JSON格式字符串,可用from_json解析为结构体后展开:
SELECT parsed_info.name, parsed_info.age, parsed_info.city FROM ( SELECT from_json(info, 'struct<name:string,age:int,city:string>') AS parsed_info FROM source_table ) t
内容的提问来源于stack exchange,提问作者Learn2Code
相关产品推荐
相关产品推荐

