【发布时间】: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