【发布时间】:2023-03-05 04:47:01
【问题描述】:
当我尝试通过sbt package创建以下包时:
import org.apache.spark.sql.SparkSession
class Log(val cip: String, val scstatus: Int) {
var src: String = cip
var status: Int = scstatus
}
object IISHttpLogs {
def main(args: Array[String]) {
val logFiles = "D:/temp/tests/wwwlogs"
val spark = SparkSession.builder.appName("LogParser").getOrCreate()
val sc = spark.sparkContext;
sc.setLogLevel("ERROR")
val logs = sc.textFile(logFiles)
import spark.implicits._
val rowDF = logs.filter(l => !l.startsWith("#"))
.map(l => l.split(" "))
.map(c => new Log(c(8), c(11).trim.toInt))
.toDF();
println(s"line count: ${rowDF.count()}")
rowDF.createOrReplaceTempView("rows")
val maxHit = spark.sql("SELECT top 1 src, count(*) FROM rows group by src order by count(*) desc")
maxHit.show()
spark.stop()
}
}
我收到以下错误:
值 toDF 不是 org.apache.spark.rdd.RDD[Log] 的成员
我尝试了几件事,例如:
- 到DFlog
- 创建一个 sql 上下文并从此 sqlContext 导入
imlicits._
我只是无法编译我的代码。
欢迎提供任何线索来覆盖此错误。
我很好地阅读了Generate a Spark StructType / Schema from a case class 并写道:
val schema =
StructType(
StructField("src", StringType, false) ::
StructField("status", IntegerType, true) :: Nil)
val rowRDD = logs.filter(l => !l.startsWith("#"))
.map(l => l.split(" "))
.map(c => Row(c(8), c(11).trim.toInt));
val rowDF = spark.sqlContext.createDataFrame(rowRDD, schema);
但是这样做我不使用Log 类。我想知道是否有办法通过使用定义的Log 类来获取DataFrame,或者官方/最佳方法是否是使用Row 类?
例如我不会写:
val rowRDD = logs.filter(l => !l.startsWith("#"))
.map(l => l.split(" "))
.map(c => new Log(c(8), c(11).trim.toInt));
val rowDF = spark.sqlContext.createDataFrame(
rowRDD,
ScalaReflection.schemaFor[Log].dataType.asInstanceOf[StructType]);
我就是不知道为什么?
【问题讨论】:
标签: scala apache-spark