【问题标题】:How to aggregate the microbatch data into a dataframe for spark structured streaming?如何将微批量数据聚合到数据帧中以进行 Spark 结构化流式传输?
【发布时间】:2020-07-14 15:40:43
【问题描述】:

用例如下。基于 spark 结构化流,摄取 kafka 数据。我们希望每个微批次的数据,每 10 秒,可以被处理并聚合到单个数据帧中,该数据帧会持续监控每个 id 的某些值的总和。下面的方法正确吗?对于每个查询,来自先前微批处理的旧数据似乎仍在“monitoring_table”中,这会导致不需要的聚合。解决此问题的最佳方法是什么?谢谢!

    val monitoring_stream = df.writeStream
                              .outputMode("append")
                              .format("memory")
                              .queryName("monitoring_table")
                              .start()

      var rounds = 0
      
      while(monitoring_stream.isActive) {
          Thread.sleep(10000)
  
          spark.sql("SELECT * from monitoring_table").show()   
          var tempDF = spark.sql("SELECT * from monitoring_table")
          var batchDF_group = tempDF.withWatermark("timestamp", "10 seconds").groupBy("id").sum("download_volume", "upload_volume").withColumnRenamed("sum(download_volume)","total_download_volume_batch").withColumnRenamed("sum(upload_volume)","total_upload_volume_batch")
          monitoring_df = monitoring_df.join(batchDF_group, monitoring_df("id") === batchDF_group("id"), "left").select(monitoring_df("id"), monitoring_df("total_download_volume"), monitoring_df("upload_volume"), monitoring_df("total_volume"), batchDF_group("total_download_volume_batch"), batchDF_group("total_upload_volume_batch")).na.fill(0)
          monitoring_df = monitoring_df.withColumn("total_upload_volume", monitoring_df("total_upload_volume")+monitoring_df("total_upload_volume_batch"))
          monitoring_df = monitoring_df.withColumn("total_download_volume", monitoring_df("total_download_volume")+monitoring_df("total_download_volume_batch"))
          monitoring_df = monitoring_df.withColumn("total_volume", monitoring_df("total_download_volume")+monitoring_df("total_upload_volume"))                                    

          
          monitoring_df.show()    
    }```

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    我已经使用 Databricks 有一段时间了,在 Streaming 中帮助我的一件事是 -

    Databricks Delta Lake

    使 Databricks Delta Lake 与此处相关的功能是 - 它使用多个流(或并发批处理作业)维护“exactly-once”处理。

    您可以浏览上述文档,如果不是databricks delta Lake,那么您可以在您的项目中包含其中的一些功能。

    【讨论】:

    • 谢谢,我去看看 Delta Take。听起来很有趣。
    猜你喜欢
    • 2018-11-08
    • 1970-01-01
    • 2019-01-14
    • 2021-10-23
    • 2019-08-09
    • 1970-01-01
    • 1970-01-01
    • 2019-07-21
    相关资源
    最近更新 更多