【问题标题】:how to use spark to analyze pv,uv,ip every 5 mins如何使用 spark 每 5 分钟分析一次 pv、uv、ip
【发布时间】:2018-01-14 09:07:21
【问题描述】:

如何每天5分钟分析一次uv、pv、ip,并存储Mysql。数据来自Kafka,格式如下:

Message sent: {"cookie":"a95f22eabc4fd4b580c011a3161a9d9d","ip":"125.119.144.252","event_time":"2017-08-07 10:50:16"}
Message sent: {"cookie":"6b67c8c700427dee7552f81f3228c927","ip":"202.109.201.181","event_time":"2017-08-07 10:50:26"}

就像 00:00-00:05 00:05--00:10 等等, 我用过:

val write=new JDBCSink()
       val query=counts.writeStream.foreach(write).outputMode("complete")
          .trigger(ProcessingTime("5 minutes"))    
          .start()

但是当我在 00:01 提交它或者它崩溃时,我怎么能确定它不会像 00:01-00:06 那样分析。

【问题讨论】:

    标签: apache-spark apache-kafka real-time bigdata


    【解决方案1】:

    使用window函数:

    query = counts.groupBy(window('event_time', '5 second')).agg()
    query.writeStream.start()
    

    【讨论】:

    • pv,uv 计算是最后一天,并且窗口不是有状态的,如果我使用这样的窗口 window($"unix_timestamp", "1 day", "5 minutes") 它也应该运行节目在 00:00 而不是 00:01
    猜你喜欢
    • 2015-07-31
    • 1970-01-01
    • 1970-01-01
    • 2013-10-19
    • 2015-09-08
    • 1970-01-01
    • 1970-01-01
    • 2020-06-23
    • 2012-06-30
    相关资源
    最近更新 更多