document.write('
Data hosted with ♥ by Pastebin.com - Download Raw - See Original
  1. # -*- coding: utf-8 -*-
  2.  
  3. """
  4. Attribute value frequency outlier detection using Spark
  5. """
  6.  
  7. from pyspark import SparkContext
  8. from operator import add
  9.  
  10. sc = SparkContext(appName="Attribute Value Frequency")
  11.  
  12. def parseDevelopmentDatasetLine(line):
  13.   """
  14.  Parses development dataset line to attributes.
  15.  
  16.  :param  line: .
  17.  :returns: list of key value pairs.
  18.  """
  19.  
  20.   retval = []
  21.   # Development dataset uses \', \' to separate
  22.   # the attributes.
  23.   parts = line.split(",")
  24.   for item in parts:
  25.     # Keys are separated from the values by \':\'
  26.     key, value = item.split(\':\')
  27.     # Lets trim the extra whitespace from the key and the value
  28.     key = key.strip()
  29.     value = value.strip()
  30.     retval.append((key, value))
  31.  
  32.   return retval
  33.  
  34. def attributeValueFrequencyScorer(line):
  35.   """
  36.  Connects each line to attribute value frequency score (up to constant factor)
  37.  
  38.  :param  line: line turned to list of (label, value) tuples
  39.  :returns: original line and attribute value frequency score.
  40.  """
  41.  
  42.   # Let\'s score each line and create a tuple (line, score)
  43.   score = 0
  44.   for item in line:
  45.     # item[0] is attribute id
  46.     # item[1] is value of the attribute
  47.  
  48.     # Note: May fail due to overflow but not likely.
  49.     score += counts[str(item[0])][str(item[1])]
  50.  
  51.   return (line, score)
  52.  
  53. if __name__ == \'__main__\':
  54.  
  55.   data = sc.textFile("s3a://<your aws account key>:<your aws account secret key>@free-style-sciencey-spark-demo/development_dataset.txt")
  56.   parsed_data = data.map(parseDevelopmentDatasetLine)
  57.  
  58.   # Let\'s count the number of times different attribute values
  59.   # are present in the dataset.
  60.  
  61.   grouped_by_attribute_value_labels = parsed_data.flatMap(lambda x: x).reduceByKey(lambda x,y: x + \' \' + y )
  62.  
  63.   keys = grouped_by_attribute_value_labels.keys().collect()
  64.  
  65.   # Count for all possible values for each attribute value
  66.   counts = {}
  67.   for key in keys:
  68.     # Let\'s turn tuples to dict
  69.     values = grouped_by_attribute_value_labels.filter(lambda x: x[0] == key).flatMap(lambda x: x[1].split()).map(lambda x: (x, 1)).reduceByKey(add).collect()
  70.     counts[str(key)] = dict(values)
  71.  
  72.   # Let\'s go through the dataset and score each item
  73.   avf_scored_data = parsed_data.map(attributeValueFrequencyScorer).sortBy(lambda x: x[1])
  74.  
  75.   outliers = avf_scored_data.take(10) # Let\'s print 10 most likely outliers
  76.   for item in outliers:
  77.     print item
  78.  
  79.   # Let\'s also put results to s3
  80.   sc.parallelize(outliers).saveAsTextFile("s3a://<your aws account key>:<your aws account secret key>@free-style-sciencey-spark-demo/avf_tuloksia.txt")
');