如何从Avro类对象组成的Scala List中获取RDD或DataSet用于测试
实现方案
前置准备
测试场景先初始化Spark上下文实例,无需连接集群:
import org.apache.spark.SparkConf import org.apache.spark.SparkContext import org.apache.spark.rdd.RDD import org.apache.spark.sql.SparkSession // 本地运行配置 val conf = new SparkConf() .setAppName("mock-test") .setMaster("local[*]") // SparkContext 用于创建RDD val sc: SparkContext = SparkContext.getOrCreate(conf) // SparkSession 用于创建DataSet val spark: SparkSession = SparkSession.builder().config(conf).getOrCreate()
转换为RDD[SomeClass]
直接调用SparkContext的parallelize方法即可,适配你的目标声明:
val someClasses: RDD[SomeClass] = sc.parallelize(inputList) // 若需要控制分区数,可传入第二个参数,示例为2个分区: // val someClasses: RDD[SomeClass] = sc.parallelize(inputList, 2)
转换为DataSet[SomeClass]
需要先定义对应Avro类的序列化编码器,分两种场景:
场景1:需要操作DataSet内部字段
引入官方spark-avro依赖后,使用Avro专用编码器,保留完整字段结构:
import org.apache.spark.sql.avro.Encoders.avro implicit val someClassEncoder = avro[SomeClass] val someClassesDs = inputList.toDS()
场景2:仅Mock传参无需操作字段
单元测试快速验证可直接用Kryo编码器,无需额外依赖:
implicit val someClassEncoder = org.apache.spark.sql.Encoders.kryo[SomeClass] val someClassesDs = inputList.toDS()
内容的提问来源于stack exchange,提问作者Vadim
相关产品推荐
相关产品推荐

