【问题标题】:spark streaming not able to use spark sql火花流无法使用火花sql
【发布时间】:2018-12-25 01:47:50
【问题描述】:

我在火花流式传输期间遇到问题。在流式传输并传递给“解析”方法后,我得到了空记录。

我的代码:

import spark.implicits._
import org.apache.spark.sql.types._
import org.apache.spark.sql.Encoders
import org.apache.spark.streaming._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.SparkSession
import spark.implicits._
import org.apache.spark.sql.types.{StructType, StructField, StringType, 
IntegerType}
import org.apache.spark.sql.functions._
import org.apache.spark.sql.SparkSession
import spark.implicits._
import org.apache.spark.sql.types.{StructType, StructField, StringType, 
IntegerType}
import org.apache.spark.SparkConf
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.storage.StorageLevel
import java.util.regex.Pattern
import java.util.regex.Matcher
import org.apache.spark.sql.hive.HiveContext;
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.sql._

val conf = new SparkConf().setAppName("streamHive").setMaster("local[*]").set("spark.driver.allowMultipleContexts", "true")

val ssc = new StreamingContext(conf, Seconds(5))    

val sc=ssc.sparkContext

val lines = ssc.textFileStream("file:///home/sadr/testHive")

case class Prices(name: String, age: String,sex: String, location: String)
val sqlContext = new org.apache.spark.sql.SQLContext(sc)

def parse (rdd : org.apache.spark.rdd.RDD[String] ) = 
{
var l = rdd.map(_.split(","))
val prices = l.map(p => Prices(p(0),p(1),p(2),p(3)))
val pricesDf = sqlContext.createDataFrame(prices)
pricesDf.registerTempTable("prices")
pricesDf.show()
var x = sqlContext.sql("select count(*) from prices")
x.show()}
lines.foreachRDD { rdd => parse(rdd)} 
lines.print()
ssc.start()

我的输入文件:

cat test1.csv

Riaz,32,M,uk
tony,23,M,india
manu,33,M,china
imart,34,F,AUS

我得到这个输出:

lines.foreachRDD { rdd => parse(rdd)}

lines.print()

ssc.start()

scala> +----+---+---+--------+
|name|age|sex|location|
+----+---+---+--------+
+----+---+---+--------+

我使用的是 Spark 2.3 版....添加 X.SHOW() 后出现以下错误

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    不确定您是否真的能够读取流。

    textFileStream 仅读取程序启动后添加到目录中的新文件,而不读取现有文件。文件已经在那里了吗? 如果是,将其从目录中删除,启动程序并再次复制文件?

    【讨论】:

    • 我在再次复制文件后运行了程序,它可以工作,但是代码中的 sql 语句 (var x = sqlContext.sql("select count(*) from prices") 没有打印任何结果。
    • x 应该是一个 DF。你可以尝试做 x.show() 而不是 println(x)
    • 嗨@Mohd Avais ...在添加 x.show() 后它不起作用..请帮助
    猜你喜欢
    • 2018-08-09
    • 2017-04-27
    • 1970-01-01
    • 2019-04-02
    • 2016-02-07
    • 2015-05-15
    • 1970-01-01
    • 1970-01-01
    • 2011-01-22
    相关资源
    最近更新 更多