Scala如何将类类型作为方法参数用于newAPIHadoopFile调用
问题解决:Scala Spark动态传入类类型调用newAPIHadoopFile
错误根源
classOf[T]是Scala编译期操作符,必须接收编译期确定的类字面量,不能传入字符串变量、普通变量动态获取类;同时你将类类型参数定义为String类型,也不符合newAPIHadoopFile的入参要求,该方法的第2-4位入参要求是Class类型对象。
实现方案
方案1:直接传入Class类型参数(最通用,类型安全)
直接修改参数类型,把三个字符串参数替换为对应的Class类型,调用时直接传类对象即可,不需要再用classOf包裹:
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.mapreduce.InputFormat import org.apache.spark.sql.SparkSession import org.apache.spark.sql.DataFrame // 泛型约束保证类型安全 case class SequenceInput[F <: InputFormat[_, _], K, V]( conf: Configuration, path: String, inputFormatClass: Class[F], keyClass: Class[K], valueClass: Class[V] ){ def read(sparkSession: SparkSession): DataFrame = { val rdd = sparkSession.sparkContext.newAPIHadoopFile( path, inputFormatClass, // 直接传Class对象即可,不需要classOf包裹 keyClass, valueClass, conf ) // 按需补充RDD转DataFrame的逻辑,示例为HBase Result解析参考 import sparkSession.implicits._ rdd.map(_._2).map(result => { val rowkey = new String(result.getRow) val colVal = new String(result.getValue("cf".getBytes, "column".getBytes)) (rowkey, colVal) }).toDF("rowkey", "column_value") } } // 调用示例 val hbaseReader = SequenceInput( conf = yourHbaseConf, path = "hdfs://path/to/hbase/file", inputFormatClass = classOf[org.apache.hadoop.mapreduce.lib.input.SequenceFileInputFormat[ImmutableBytesWritable, Result]], keyClass = classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable], valueClass = classOf[org.apache.hadoop.hbase.client.Result] ) val resultDF = hbaseReader.read(spark)
方案2:类名需用字符串传递的场景(比如类名从配置文件读取)
如果必须接收字符串形式的类名,可通过反射加载得到Class对象:
case class SequenceInput( conf: Configuration, path: String, inputFormatClassName: String, keyClassName: String, valueClassName: String ){ def read(sparkSession: SparkSession): DataFrame = { // 反射加载类并做类型转换 val inputFormatClass = Class.forName(inputFormatClassName).asInstanceOf[Class[org.apache.hadoop.mapreduce.InputFormat[_, _]]] val keyClass = Class.forName(keyClassName) val valueClass = Class.forName(valueClassName) val rdd = sparkSession.sparkContext.newAPIHadoopFile( path, inputFormatClass, keyClass, valueClass, conf ) // 同上补充RDD转DataFrame逻辑即可 } }
内容的提问来源于stack exchange,提问作者Ayush Goyal
相关产品推荐
相关产品推荐

