如何将JSON字符串转为Scala DataFrame?求toDS替代方案
替代方案:将JSON字符串转换为Spark DataFrame
原代码中Seq(abc).toDS()依赖SparkSession的隐式编码器,打包环境中可能因隐式作用域失效、依赖缺失等问题报错,以下是几种无需依赖隐式转换的可靠实现方式:
方法一:使用RDD[String]作为输入
直接将JSON字符串包装成RDD传入spark.read.json,避开DS的隐式转换要求:val spark = SparkSession.builder().getOrCreate() val abc="""[{"orders":{"order_id":{"path":"orderid","type":"string"},"customer_id":{"path":"customers.customerId","type":"string"},"offer_id":{"path":"Offers.Offerid","type":"string"}},"products":{"product_id":{"path":"product_id","type":"string"},"product_name":{"path":"products.productname","type":"string"}}}]""" val df = spark.read.json(spark.sparkContext.parallelize(Seq(abc))) df.show()方法二:使用StringReader直接读取
利用Spark支持从输入流读取JSON的特性,直接传入StringReader对象:import java.io.StringReader val spark = SparkSession.builder().getOrCreate() val abc="""[{"orders":{"order_id":{"path":"orderid","type":"string"},"customer_id":{"path":"customers.customerId","type":"string"},"offer_id":{"path":"Offers.Offerid","type":"string"}},"products":{"product_id":{"path":"product_id","type":"string"},"product_name":{"path":"products.productname","type":"string"}}}]""" val df = spark.read.json(new StringReader(abc)) df.show()方法三:手动解析JSON并构建DataFrame(复杂场景适用)
若需要更精细的结构控制,可通过JSON解析库先解析字符串为Row对象,再指定Schema创建DataFrame:import com.fasterxml.jackson.databind.ObjectMapper import org.apache.spark.sql.types._ import org.apache.spark.sql.Row import scala.collection.JavaConverters._ val spark = SparkSession.builder().getOrCreate() val abc="""[{"orders":{"order_id":{"path":"orderid","type":"string"},"customer_id":{"path":"customers.customerId","type":"string"},"offer_id":{"path":"Offers.Offerid","type":"string"}},"products":{"product_id":{"path":"product_id","type":"string"},"product_name":{"path":"products.productname","type":"string"}}}]""" // 定义与JSON匹配的Schema val fieldSchema = StructType(Seq( StructField("path", StringType), StructField("type", StringType) )) val schema = StructType(Seq( StructField("orders", StructType(Seq( StructField("order_id", fieldSchema), StructField("customer_id", fieldSchema), StructField("offer_id", fieldSchema) ))), StructField("products", StructType(Seq( StructField("product_id", fieldSchema), StructField("product_name", fieldSchema) ))) )) // 解析JSON为Row集合 val mapper = new ObjectMapper() val jsonNodes = mapper.readTree(abc) val rows = jsonNodes.elements().asScala.map(node => { val orderNode = node.get("orders") val productNode = node.get("products") Row( Row( Row(orderNode.get("order_id").get("path").asText(), orderNode.get("order_id").get("type").asText()), Row(orderNode.get("customer_id").get("path").asText(), orderNode.get("customer_id").get("type").asText()), Row(orderNode.get("offer_id").get("path").asText(), orderNode.get("offer_id").get("type").asText()) ), Row( Row(productNode.get("product_id").get("path").asText(), productNode.get("product_id").get("type").asText()), Row(productNode.get("product_name").get("path").asText(), productNode.get("product_name").get("type").asText()) ) ) }).toList // 创建DataFrame val df = spark.createDataFrame(spark.sparkContext.parallelize(rows), schema) df.show()
内容的提问来源于stack exchange,提问作者Shankar Panda
相关产品推荐
相关产品推荐

