【问题标题】:Check logs with Spark使用 Spark 检查日志
【发布时间】:2016-12-29 19:57:04
【问题描述】:

我是 Spark 的新手,我正在尝试开发一个 python 脚本来读取带有一些日志的 csv 文件:

userId,timestamp,ip,event
13,2016-12-29 16:53:44,86.20.90.121,login
43,2016-12-29 16:53:44,106.9.38.79,login
66,2016-12-29 16:53:44,204.102.78.108,logoff
101,2016-12-29 16:53:44,14.139.102.226,login
91,2016-12-29 16:53:44,23.195.2.174,logoff

并检查用户是否有一些奇怪的行为,例如他是否连续两次“登录”而没有“注销”。我已将 csv 作为 Spark 数据帧加载,我想比较单个用户的日志行,按时间戳排序,并检查两个连续事件是否属于同一类型(登录 - 登录,注销 - 注销)。我正在寻找以“map-reduce”方式进行操作,但目前我无法弄清楚如何使用比较连续行的 reduce 函数。 我写的代码可以工作,但是性能很差。

sc = SparkContext("local","Data Check")
sqlContext = SQLContext(sc)

LOG_FILE_PATH = "hdfs://quickstart.cloudera:8020/user/cloudera/flume/events/*"
RESULTS_FILE_PATH = "hdfs://quickstart.cloudera:8020/user/cloudera/spark/script_results/prova/bad_users.csv"
N_USERS = 10*1000

dataFrame = sqlContext.read.format("com.databricks.spark.csv").load(LOG_FILE_PATH)
dataFrame = dataFrame.selectExpr("C0 as userID","C1 as timestamp","C2 as ip","C3 as event")

wrongUsers = []

for i in range(0,N_USERS):

    userDataFrame = dataFrame.where(dataFrame['userId'] == i)
    userDataFrame = userDataFrame.sort('timestamp')

    prevEvent = ''

    for row in userDataFrame.rdd.collect():

        currEvent = row[3]
        if(prevEvent == currEvent):
            wrongUsers.append(row[0])

        prevEvent = currEvent

badUsers = sqlContext.createDataFrame(wrongUsers)
badUsers.write.format("com.databricks.spark.csv").save(RESULTS_FILE_PATH)

【问题讨论】:

    标签: python apache-spark mapreduce pyspark spark-dataframe


    【解决方案1】:

    首先(不相关但仍然),请确保每个用户的条目数不是很大,因为for row in userDataFrame.rdd.collect(): 中的collect 很危险。

    其次,使用经典Python不需要离开DataFrame区,只要坚持使用Spark即可。

    现在,你的问题。它基本上是“对于每一行我想从上一行中了解一些东西”:属于Window 函数的概念,确切地说是lag 函数。这里有两篇关于 Spark 中 Window 函数的有趣文章:一篇来自 Databricks 的 Python 代码,另一篇来自 Xinh 的(我认为更容易理解的)Scala 示例。

    我在 Scala 中有一个解决方案,但我想你会在 Python 中翻译它:

    import org.apache.spark.sql.expressions.Window
    import org.apache.spark.sql.functions.lag
    
    import sqlContext.implicits._
    
    val LOG_FILE_PATH = "hdfs://quickstart.cloudera:8020/user/cloudera/flume/events/*"
    val RESULTS_FILE_PATH = "hdfs://quickstart.cloudera:8020/user/cloudera/spark/script_results/prova/bad_users.csv"
    
    val data = sqlContext
      .read
      .format("com.databricks.spark.csv")
      .option("inferSchema", "true")
      .option("header", "true") // use the header from your csv
      .load(LOG_FILE_PATH)
    
    val wSpec = Window.partitionBy("userId").orderBy("timestamp")
    
    val badUsers = data
      .withColumn("previousEvent", lag($"event", 1).over(wSpec))
      .filter($"previousEvent" === $"event")
      .select("userId")
      .distinct
    
    badUsers.write.format("com.databricks.spark.csv").save(RESULTS_FILE_PATH)
    

    基本上,您只需从上一行检索值并将其与当前行上的值进行比较,如果匹配是错误行为并且您保留userId。对于每个userId 的行“块”中的第一行,之前的值将是null:当与当前值比较时,布尔表达式将为false,所以这里没问题。

    【讨论】:

    • 非常感谢,这正是我正在寻找的解决方案。也可以将函数传递给过滤方法吗?因为我还得用 ip 做一些事情
    • 是的,您当然可以在过滤器中添加一些内容。请记住,如果您无法通过简单的比较或 Spark 的基于列的functions 表达您希望函数执行的操作,则需要使用UDFs,请参阅page 了解更多信息。
    猜你喜欢
    • 2020-06-12
    • 2018-11-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-02
    • 1970-01-01
    • 2015-04-07
    • 1970-01-01
    相关资源
    最近更新 更多