【问题标题】:Bigram frequency of query logs events with Apache Spark使用 Apache Spark 的查询日志事件的 Bigram 频率
【发布时间】:2015-04-07 04:57:30
【问题描述】:

我想研究从搜索引擎查询日志中提取的会话中的用户操作。我定义了前两种动作:查询和点击。

sealed trait Action{}
case class Query(val input:String) extends Action
case class Click(val link:String)  extends Action

假设查询日志中的第一个操作由以下时间戳(以毫秒为单位)给出:

val t0 = 1417444964686L // 2014-12-01 15:42:44

让我们定义一个与会话 ID 关联的按时间排序的操作语料库。

val query_log:Array[(String, (Action, Long))] = Array (
("session1",(Query("query1"),t0)), 
("session1",(Click("link1") ,t0+1000)), 
("session1",(Click("link2") ,t0+2000)), 
("session1",(Query("query2"),t0+3000)), 
("session1",(Click("link3") ,t0+4000)), 
("session2",(Query("query3"),t0+5000)), 
("session2",(Click("link4") ,t0+6000)), 
("session2",(Query("query4"),t0+7000)), 
("session2",(Query("query5"),t0+8000)),
("session2",(Click("link5") ,t0+9000)),
("session2",(Click("link6") ,t0+10000)),
("session3",(Query("query6"),t0+11000))
)

我们为此 quey_log 创建一个 RDD:

import org.apache.spark.rdd.RDD
var logs:RDD[(String, (Action, Long))] = sc.makeRDD(query_log)

日志然后按会话 ID 分组

val sessions_groups:RDD[(String, Iterable[(Action, Long)])] = logs.groupByKey().cache()

现在,我们要研究会话中的动作共现,例如,会话中的重写次数。然后,我们定义了 Cooccurrences 类,该类将从会话操作中初始化。

case class Cooccurrences(
  var numQueriesWithClicks:Int = 0,
  var numQueries:Int = 0,
  var numRewritings:Int = 0,
  var numQueriesBeforeClicks:Int = 0
 ) {
 // The cooccurrence object is initialized from a list of timestamped action in order to catch a session group
  def initFromActions(actions:Iterable[(Action, Long)]) = {
    // 30 seconds is the maximal time (in milliseconds) between two  queries (q1, q2) to consider q2 is a rewririting of q1
    var thirtySeconds = 30000 
    var hasClicked = false 
    var hasRewritten = false
    // int the observed action sequence, we extract consecutives (sliding(2)) actions sorted by timestamps
    // for each bigram in the sequence we want to count and modify the cooccurrence object
    actions.toSeq.sortBy(_._2).sliding(2).foreach{ 
      // case Seq(l0) => // session with only one Action 
      case Seq((e1:Click, t0)) => { // click without any query
        numQueries = 0        
      }
      case Seq((e1:Query, t0)) => { // query without any click
        numQueries = 1        
        numQueriesBeforeClicks = 1
      }
      // case Seq(l0, l1) => // session with at least two Actions
      case Seq((e1:Click, t0), (e2:Query, t1)) => { // a click followed by a query
        if(! hasClicked)
          numQueriesBeforeClicks = numQueries
        hasClicked = true
        }
      case Seq((e1:Click, t0), (e2:Click, t1)) => { //two consecutives clics 
        if(! hasClicked)
          numQueriesBeforeClicks = numQueries
        hasClicked = true
      }
      case Seq((e1:Query, t0), (e2:Click, t1)) => { // a query followed by a click
        numQueries += 1
        if(! hasClicked)
          numQueriesBeforeClicks = numQueries
        hasClicked = true
        numQueriesWithClicks +=1
      }
      case Seq((e1:Query, t0), (e2:Query, t1)) => { // two consecutives queries
        val dt = t1 - t0
        numQueries += 1
        if(dt < thirtySeconds && e1.input != e2.input){
          hasRewritten = true
          numRewritings += 1
       }
      }
    }
  }

}

现在,让我们尝试计算每个会话的共现 RDD:

val session_cooc_stats:RDD[Cooccurrences] = sessions_groups.map{ 
  case (sessionId, actions) => {
   var coocs  = Cooccurrences()
   coocs.initFromActions(actions)
   coocs
  }
 }

不幸的是,它引发了以下 MatchError

scala> session_cooc_stats.take(2)

15/02/06 22:50:08 ERROR Executor: Exception in task 0.0 in stage 1.0 (TID 4) scala.MatchError: List((Query(query3),1417444969686), (Click(link4),1417444970686)) (of class scala.collection.immutable.$colon$colon) at $line25.$read$$iwC$$iwC$Cooccurrences$$anonfun$initFromActions$2.apply(<console>:29)
  at $line25.$read$$iwC$$iwC$Cooccurrences$$anonfun$initFromActions$2.apply(<console>:29)
  at scala.collection.Iterator$class.foreach(Iterator.scala:727)
  at scala.collection.AbstractIterator.foreach(Iterator.scala:1157)
  at $line25.$read$$iwC$$iwC$Cooccurrences.initFromActions(<console>:29)
  at $line28.$read$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:31)
  at $line28.$read$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:28)
  at scala.collection.Iterator$$anon$11.next(Iterator.scala:328)
  at scala.collection.Iterator$$anon$10.next(Iterator.scala:312)
  at scala.collection.Iterator$class.foreach(Iterator.scala:727)
  at scala.collection.AbstractIterator.foreach(Iterator.scala:1157)
  at scala.collection.generic.Growable$class.$plus$plus$eq(Growable.scala:48)
  at scala.collection.mutable.ArrayBuffer.$plus$plus$eq(ArrayBuffer.scala:103)
  at scala.collection.mutable.ArrayBuffer.$plus$plus$eq(ArrayBuffer.scala:47)
  at scala.collection.TraversableOnce$class.to(TraversableOnce.scala:273)
  at scala.collection.AbstractIterator.to(Iterator.scala:1157)
  at scala.collection.TraversableOnce$class.toBuffer(TraversableOnce.scala:265)
  at scala.collection.AbstractIterator.toBuffer(Iterator.scala:1157)
  at scala.collection.TraversableOnce$class.toArray(TraversableOnce.scala:252)
  at scala.collection.AbstractIterator.toArray(Iterator.scala:1157)
  at org.apache.spark.rdd.RDD$$anonfun$26.apply(RDD.scala:1081)
  at org.apache.spark.rdd.RDD$$anonfun$26.apply(RDD.scala:1081)
  at org.apache.spark.SparkContext$$anonfun$runJob$4.apply(SparkContext.scala:1314)
  at org.apache.spark.SparkContext$$anonfun$runJob$4.apply(SparkContext.scala:1314)
  at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:61)
  at org.apache.spark.scheduler.Task.run(Task.scala:56)
  at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:196)
  at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
  at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
  at java.lang.Thread.run(Thread.java:745)
15/02/06 22:50:08 WARN TaskSetManager: Lost task 0.0 in stage 1.0 (TID 4, localhost): scala.MatchError: List((Query(query3),1417444969686), (Click(link4),1417444970686)) (of class scala.collection.immutable.$colon$colon)
  at $line25.$read$$iwC$$iwC$Cooccurrences$$anonfun$initFromActions$2.apply(<console>:29)
  at $line25.$read$$iwC$$iwC$Cooccurrences$$anonfun$initFromActions$2.apply(<console>:29)
 ...

如果我建立自己的行动列表相当于 session_cooc_stats RDD 中的第一组

val actions:Iterable[(Action, Long)] = Array(
(Query("query1"),t0),
(Click("link1") ,t0+1000),
(Click("link2") ,t0+2000),
(Query("query2"),t0+3000),
(Click("link3") ,t0+4000)
)

我得到了预期的结果

var c = Cooccurrences()
c.initFromActions(actions)
// c == Cooccurrences(2,2,0,1)

当我从 RDD 构建 Cooccurrence 对象时,似乎出现了问题。 它似乎与使用 groupByKey() 构建的 CompactBuffer 相关联。
缺什么 ?

我是 Spark 和 Scala 的新手。 提前感谢您的帮助。

托马斯

【问题讨论】:

  • 我刚插上它,它对我有用...
  • 你真的可以在不做任何更改的情况下执行Spark指令session_cooc_stats.take(2)吗?我仍然在 List((Query(query3),1417444969686), (Click(link4),1417444970686)) 上得到一个 scala.MatchError

标签: scala session pattern-matching apache-spark sequence


【解决方案1】:

按照您的建议,我用 IntelliJ 重写了代码并为 main 函数创建了一个伴随对象。 令人惊讶的是,代码可以编译(使用 sbt)并完美运行。

但是,我真的不明白为什么编译的代码可以运行,而它不能与 spark-shell 一起使用。

感谢您的回答!

【讨论】:

    【解决方案2】:

    我将您的代码设置在 IntelliJ 上。

    为 Action、Query、Click 和 Coocurence 创建一个类。

    你的代码放在一个主目录上。

    val sessions_groups:RDD[(String, Iterable[(Action, Long)])] = logs.groupByKey().cache()
    
      val session_cooc_stats:RDD[Cooccurrences] = sessions_groups.map{
        case (sessionId, actions) => {
          val coocs  = Cooccurrences()
          coocs.initFromActions(actions)
          coocs
        }
      }
      session_cooc_stats.take(2).foreach(println(_))
    

    刚刚修改的 var coocs > val coocs

    我猜是重点。

    同时出现(0,1,0,1)

    同时出现(2,3,1,1)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多