【问题标题】:spark streaming and spark sql considerationsspark流和spark sql注意事项
【发布时间】:2018-10-13 07:42:36
【问题描述】:

我正在使用 spark 流 (scala) 并在每 20 分钟后通过 kafka 接收客户呼叫呼叫中心的记录。这些记录在 rdd 和更高版本的数据帧中转换以利用 spark sql。我有一个业务用例,我想识别在过去两个小时内打电话超过 3 次的所有客户。

最好的方法是什么?我是否应该继续在配置单元表中插入每批收到的所有记录并运行单独的脚本来继续查询谁在过去两个小时内进行了 3 次调用,或者还有另一个更好地使用 spark 的内存功能?

谢谢。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql spark-dataframe spark-streaming


    【解决方案1】:

    对于这种用例,您可以使用 spark 获得结果(不需要 hive)。您必须有一些客户唯一 ID,因此您可以准备一些查询,例如

    ROW_NUMBER() OVER(PARTITION BY cust_id ORDER BY time DESC) as call_count
    

    您必须使用call_count=3 对最近 2 小时的数据应用过滤器,以便获得预期的结果。然后,您可以将此 spark 脚本设置为 crontab 或任何其他自动运行方法。

    【讨论】:

    • 谢谢@Sahil。但我每 20 分钟分批接收一次数据。我得到 RDD 并将它们转换为数据帧,然后是 tempview。按照这个逻辑,我的 tempview 将只包含最后一个 RDD。如何确保我保存了不同批次的所有 RDD 的结果,然后应用 row_number() 的逻辑下面是我的示例
    • AirDRStream.foreachRDD(foreachFunc = rdd => { System.out.println("--- New RDD with " + rdd.count() + " records"); val sqlContext = SparkSession .builder () .appName("Spark SQL 基本示例") .getOrCreate() import sqlContext.implicits._ rdd.toDF().createOrReplaceTempView("AIR") val FilteredDR = sqlContext.sql("select refillProfileID, count(*) from AIR group by refillProfileID") FilteredDR.show()
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-03
    • 2014-10-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多