【问题标题】:How to auto update %spark.sql result in zeppelin for structured streaming query如何在 zeppelin 中自动更新 %spark.sql 结果以进行结构化流式查询
【发布时间】:2017-12-18 09:45:29
【问题描述】:

我正在(spark 2.1.0 和 zeppelin 0.7)中针对来自 kafka 的数据运行结构化流式传输,并且我正在尝试使用 spark.sql 可视化流式传输结果

如下:

%spark2
val spark = SparkSession
  .builder()
  .appName("Spark structured streaming Kafka example")
  .master("yarn")
  .getOrCreate()
val inputstream = spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "n11.hdp.com:6667,n12.hdp.com:6667,n13.hdp.com:6667 ,n10.hdp.com:6667, n9.hdp.com:6667")
    .option("subscribe", "st")
    .load()


val stream = inputstream.selectExpr("CAST( value AS STRING)").as[(String)].select(
             expr("(split(value, ','))[0]").cast("string").as("pre_post_paid"),
             expr("(split(value, ','))[1]").cast("double").as("DataUpload"),
             expr("(split(value, ','))[2]").cast("double").as("DataDowndownload"))
           .filter("DataUpload is not null and DataDowndownload is not null")
          .groupBy("pre_post_paid").agg(sum("DataUpload") + sum("DataDowndownload") as "size")
val query = stream.writeStream
.format("memory")
.outputMode("complete")
.queryName("test")
.start()

在它运行之后,我在“测试”上查询如下:

%sql
select *
from test

它仅在我手动运行时更新,我的问题是如何在处理新数据时更新它(流式可视化),例如:

Insights Without Tradeoffs: Using Structured Streaming in Apache Spark

【问题讨论】:

  • 我在看类似的东西。你有想过这个吗?

标签: apache-spark-sql spark-streaming visualization apache-zeppelin


【解决方案1】:

换行

"%sql    
select *
from test"

%spark
spark.table("test").show()

【讨论】:

    猜你喜欢
    • 2021-10-22
    • 2020-01-30
    • 2021-02-06
    • 2018-08-11
    • 1970-01-01
    • 2018-08-10
    • 2019-05-01
    • 2016-07-01
    • 1970-01-01
    相关资源
    最近更新 更多