【发布时间】:2019-06-28 04:26:43
【问题描述】:
我正在尝试创建从 Kafka 到 Spark 的结构化流,这是一个 json 字符串。现在想要将 json 解析为特定列,然后以最佳速度将数据帧保存到 cassandra 表中。使用 Spark 2.4 和 cassandra 2.11 (Apache) 而不是 DSE。
我尝试创建一个 Direct Stream,它提供案例类的 DStream,我在 DStream 上使用 foreachRDD 将其保存到 Cassandra 中,但这每 6-7 天后就会挂起。因此尝试流式传输直接提供数据帧并可以保存到 Cassandra。
val conf = new SparkConf()
.setMaster("local[3]")
.setAppName("Fleet Live Data")
.set("spark.cassandra.connection.host", "ip")
.set("spark.cassandra.connection.keep_alive_ms", "20000")
.set("spark.cassandra.auth.username", "user")
.set("spark.cassandra.auth.password", "pass")
.set("spark.streaming.stopGracefullyOnShutdown", "true")
.set("spark.executor.memory", "2g")
.set("spark.driver.memory", "2g")
.set("spark.submit.deployMode", "cluster")
.set("spark.executor.instances", "4")
.set("spark.executor.cores", "2")
.set("spark.cores.max", "9")
.set("spark.driver.cores", "9")
.set("spark.speculation", "true")
.set("spark.locality.wait", "2s")
val spark = SparkSession
.builder
.appName("Fleet Live Data")
.config(conf)
.getOrCreate()
println("Spark Session Config Done")
val sc = SparkContext.getOrCreate(conf)
sc.setLogLevel("ERROR")
val ssc = new StreamingContext(sc, Seconds(10))
val sqlContext = new SQLContext(sc)
val topics = Map("livefleet" -> 1)
import spark.implicits._
implicit val formats = DefaultFormats
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "brokerIP:port")
.option("subscribe", "livefleet")
.load()
val collection = df.selectExpr("CAST(value AS STRING)").map(f => parse(f.toString()).extract[liveevent])
val query = collection.writeStream
.option("checkpointLocation", "/tmp/check_point/")
.format("kafka")
.format("org.apache.spark.sql.cassandra")
.option("keyspace", "trackfleet_db")
.option("table", "locationinfotemp1")
.outputMode(OutputMode.Update)
.start()
query.awaitTermination()
预期是将数据帧保存到 cassandra。但是收到此错误:-
线程 "main" org.apache.spark.sql.AnalysisException 中的异常:必须使用 writeStream.start() 执行带有流式源的查询
【问题讨论】:
-
你看过 Kafka Connect 吗?这是 Apache Kafka 的一部分,是一种将数据从 Kafka 主题流式传输到目标数据存储(例如 Cassandra)的好方法。
-
提示:
.format("kafka").format("org.apache.spark.sql.cassandra")不正确 -
你有没有在代码末尾调用
writeStream.start()? -
@cricket_007 - 我知道它不正确,但我实际上正在寻找解决方案,如果我删除 .format("org.apache.spark.sql.cassandra") 这个应该是正确的部分,然后它可以工作,但在这种情况下,它会开始在控制台上显示并且不会保存到 cassandra。
-
@Pinnacle 是的,Kafka Connect 可以分布式运行。如果您想使用 Spark,那很好,我只是在检查您是否知道可能更适合的替代工具。
标签: scala apache-kafka spark-streaming spark-structured-streaming spark-cassandra-connector