Skip to article frontmatterSkip to article content
Site not loading correctly?

This may be due to an incorrect BASE_URL configuration. See the MyST Documentation for reference.

import json
import re
class DisplayRDD:
        def __init__(self, rdd):
                self.rdd = rdd

        def _repr_html_(self):                                  
                x = self.rdd.mapPartitionsWithIndex(lambda i, x: [(i, [y for y in x])])
                l = x.collect()
                s = "<table><tr>{}</tr><tr><td>".format("".join(["<th>Partition {}".format(str(j)) for (j, r) in l]))
                s += '</td><td valign="bottom" halignt="left">'.join(["<ul><li>{}</ul>".format("<li>".join([str(rr) for rr in r])) for (j, r) in l])
                s += "</td></table>"
                return s
## Load data into RDDs
playRDD = sc.textFile("datafiles/play.txt")
logsRDD = sc.textFile("Assignment-5-Autograder/bigdatafiles/NASA_logs_sample.txt")
amazonInputRDD = sc.textFile("datafiles/amazon-ratings.txt")
nobelRDD = sc.textFile("datafiles/prize.json")

## The following converts the amazonInputRDD into 2-tuples with integers
amazonBipartiteRDD = amazonInputRDD.map(lambda x: x.split(" ")).map(lambda x: (x[0], x[1])).distinct()

# A hack to avoid having to pass 'sc' around
dummyrdd = None
def setDefaultAnswer(rdd): 
    global dummyrdd
    dummyrdd = rdd
setDefaultAnswer(sc.parallelize([0]))
def task1(amazonInputRDD):
        return dummyrdd
task1_result = task1(amazonInputRDD)
for x in task1_result.takeOrdered(10):
    print(x)
0
def task2(amazonInputRDD):
        return dummyrdd
task2_result = task2(amazonInputRDD)
for x in task2_result.takeOrdered(10):
    print(x)
0
def task3(amazonInputRDD):
        return dummyrdd
task3_result = task3(amazonInputRDD)
for x in task3_result.takeOrdered(10):
    print(x)      
0
def task4(logsRDD):
        return dummyrdd
task4_result = task4(logsRDD)
for x in task4_result.takeOrdered(10):
    print(x)
0
def task5_flatmap(x):
        return []

task5_result = playRDD.flatMap(task5_flatmap).distinct()
print(task5_result.takeOrdered(100))
[]
def task6(playRDD):
        return dummyrdd

task6_result = task6(playRDD)
for x in task6_result.takeOrdered(10):
    print(x)
0
def task7_flatmap(x):
        return []

task7_result = nobelRDD.map(json.loads).flatMap(task7_flatmap).distinct()
print(task7_result.takeOrdered(10))
[]
def task8(nobelRDD):
        return dummyrdd

task8_result = task8(nobelRDD)
for x in task8_result.takeOrdered(10):
        print(x)
0
def task9(logsRDD, l):
        return dummyrdd
def task91(logsRDD, l):
        def extractHost(logline):
                match = re.search('^(\S+) ', logline)
                return match.group(1) if match is not None else None
        r1 = logsRDD.map(lambda s: (extractHost(s), [d in s for d in l]))
        #r2 = r1.reduceByKey(lambda x1, x2: (x1[0] or x2[0], x1[1] or x2[1]))
        r2 = r1.reduceByKey(lambda x1, x2: tuple([x1[i] or x2[i] for i in range(0, len(l))]))
        return r2.filter(lambda x: all(v for v in x[1])).map(lambda x: x[0])


def task92(logsRDD, l):
    return logsRDD.map(lambda x: ([word for word in x.split(" ") if word != ''][0], re.search(r"\d{2}/[a-zA-Z]{3}/\d{4}", x).group()))

        #return logsRDD.map(lambda x: ([word for word in x.split(" ") if word != ''][0], re.search(r"\d{2}/[a-zA-Z]{3}/\d{4}", x).group())).groupByKey().map(lambda x: x[0] if set(dict.fromkeys(x[1]))==set(l) else '').filter(lambda x: x!='').distinct()

def task93(logsRDD, l):
        def helper1(line):
                array_of_words = line.split()
                if (array_of_words[3][1: 12] in l):
                        return (array_of_words[0], (array_of_words[3][1: 12]))
                        


        def helper2(line):
                in_all = True
                for date in l:
                        if(date not in list(line[1])):
                                in_all = False
                return in_all   

        def helper3(line):
                return (line[0], list(line[1]))
                
        
                
        logs_rdd_1 = logsRDD.map(helper1).distinct()
        print(logs_rdd_1.count())
        logs_rdd_2 = logs_rdd_1.groupByKey()
        print(logs_rdd_2.count())
        logs_rdd_3 = logs_rdd_2.filter((helper2))
        print(logs_rdd_3.count())
        logs_rdd_4 = logs_rdd_3.map(lambda x: (x[0]))
        print(logs_rdd_4.count())
        return logs_rdd_4

    
task9_result = task93(logsRDD, ['02/Jul/1995', '03/Jul/1995', '04/Jul/1995', '05/Jul/1995', '06/Jul/1995'])
print(task9_result.count())
for x in task9_result.takeOrdered(10):
    print(x)
72400
---------------------------------------------------------------------------
Py4JJavaError                             Traceback (most recent call last)
<ipython-input-18-47210123d3e9> in <module>
     47 
     48 
---> 49 task9_result = task93(logsRDD, ['02/Jul/1995', '03/Jul/1995', '04/Jul/1995', '05/Jul/1995', '06/Jul/1995'])
     50 print(task9_result.count())
     51 for x in task9_result.takeOrdered(10):

<ipython-input-18-47210123d3e9> in task93(logsRDD, l)
     39         print(logs_rdd_1.count())
     40         logs_rdd_2 = logs_rdd_1.groupByKey()
---> 41         print(logs_rdd_2.count())
     42         logs_rdd_3 = logs_rdd_2.filter((helper2))
     43         print(logs_rdd_3.count())

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py in count(self)
   1235         3
   1236         """
-> 1237         return self.mapPartitions(lambda i: [sum(1 for _ in i)]).sum()
   1238 
   1239     def stats(self):

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py in sum(self)
   1224         6.0
   1225         """
-> 1226         return self.mapPartitions(lambda x: [sum(x)]).fold(0, operator.add)
   1227 
   1228     def count(self):

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py in fold(self, zeroValue, op)
   1078         # zeroValue provided to each partition is unique from the one provided
   1079         # to the final reduce call
-> 1080         vals = self.mapPartitions(func).collect()
   1081         return reduce(op, vals, zeroValue)
   1082 

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py in collect(self)
    948         """
    949         with SCCallSiteSync(self.context) as css:
--> 950             sock_info = self.ctx._jvm.PythonRDD.collectAndServe(self._jrdd.rdd())
    951         return list(_load_from_socket(sock_info, self._jrdd_deserializer))
    952 

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/py4j-0.10.9.2-src.zip/py4j/java_gateway.py in __call__(self, *args)
   1307 
   1308         answer = self.gateway_client.send_command(command)
-> 1309         return_value = get_return_value(
   1310             answer, self.gateway_client, self.target_id, self.name)
   1311 

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/sql/utils.py in deco(*a, **kw)
    109     def deco(*a, **kw):
    110         try:
--> 111             return f(*a, **kw)
    112         except py4j.protocol.Py4JJavaError as e:
    113             converted = convert_exception(e.java_exception)

/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/py4j-0.10.9.2-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name)
    324             value = OUTPUT_CONVERTER[type](answer[2:], gateway_client)
    325             if answer[1] == REFERENCE_TYPE:
--> 326                 raise Py4JJavaError(
    327                     "An error occurred while calling {0}{1}{2}.\n".
    328                     format(target_id, ".", name), value)

Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.collectAndServe.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 30.0 failed 1 times, most recent failure: Lost task 0.0 in stage 30.0 (TID 52) (8d38a7c37afd executor driver): org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/worker.py", line 619, in main
    process()
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/worker.py", line 609, in process
    out_iter = func(split_index, iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2918, in pipeline_func
    return func(split, prev_func(split, iterator))
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2918, in pipeline_func
    return func(split, prev_func(split, iterator))
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 417, in func
    return f(iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2236, in combine
    merger.mergeValues(iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/shuffle.py", line 240, in mergeValues
    for k, v in iterator:
TypeError: cannot unpack non-iterable NoneType object

	at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.handlePythonException(PythonRunner.scala:545)
	at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:703)
	at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:685)
	at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:498)
	at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37)
	at scala.collection.Iterator$GroupedIterator.fill(Iterator.scala:1211)
	at scala.collection.Iterator$GroupedIterator.hasNext(Iterator.scala:1217)
	at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
	at org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.write(BypassMergeSortShuffleWriter.java:140)
	at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59)
	at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
	at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52)
	at org.apache.spark.scheduler.Task.run(Task.scala:131)
	at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:506)
	at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1462)
	at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:509)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
	at java.lang.Thread.run(Thread.java:748)

Driver stacktrace:
	at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2403)
	at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2352)
	at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2351)
	at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
	at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
	at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
	at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2351)
	at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1109)
	at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1109)
	at scala.Option.foreach(Option.scala:407)
	at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1109)
	at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2591)
	at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2533)
	at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2522)
	at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
	at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:898)
	at org.apache.spark.SparkContext.runJob(SparkContext.scala:2214)
	at org.apache.spark.SparkContext.runJob(SparkContext.scala:2235)
	at org.apache.spark.SparkContext.runJob(SparkContext.scala:2254)
	at org.apache.spark.SparkContext.runJob(SparkContext.scala:2279)
	at org.apache.spark.rdd.RDD.$anonfun$collect$1(RDD.scala:1030)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
	at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
	at org.apache.spark.rdd.RDD.withScope(RDD.scala:414)
	at org.apache.spark.rdd.RDD.collect(RDD.scala:1029)
	at org.apache.spark.api.python.PythonRDD$.collectAndServe(PythonRDD.scala:180)
	at org.apache.spark.api.python.PythonRDD.collectAndServe(PythonRDD.scala)
	at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
	at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
	at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.lang.reflect.Method.invoke(Method.java:498)
	at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
	at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
	at py4j.Gateway.invoke(Gateway.java:282)
	at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
	at py4j.commands.CallCommand.execute(CallCommand.java:79)
	at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
	at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
	at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.spark.api.python.PythonException: Traceback (most recent call last):
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/worker.py", line 619, in main
    process()
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/worker.py", line 609, in process
    out_iter = func(split_index, iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2918, in pipeline_func
    return func(split, prev_func(split, iterator))
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2918, in pipeline_func
    return func(split, prev_func(split, iterator))
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 417, in func
    return f(iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/pyspark/rdd.py", line 2236, in combine
    merger.mergeValues(iterator)
  File "/data/Assignment-5/spark-3.2.0-bin-hadoop3.2/python/lib/pyspark.zip/pyspark/shuffle.py", line 240, in mergeValues
    for k, v in iterator:
TypeError: cannot unpack non-iterable NoneType object

	at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.handlePythonException(PythonRunner.scala:545)
	at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:703)
	at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:685)
	at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:498)
	at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37)
	at scala.collection.Iterator$GroupedIterator.fill(Iterator.scala:1211)
	at scala.collection.Iterator$GroupedIterator.hasNext(Iterator.scala:1217)
	at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
	at org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.write(BypassMergeSortShuffleWriter.java:140)
	at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59)
	at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
	at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52)
	at org.apache.spark.scheduler.Task.run(Task.scala:131)
	at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:506)
	at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1462)
	at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:509)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
	... 1 more
def task10(bipartiteGraphRDD):
        return dummyrdd
task10_result = task10(amazonBipartiteRDD)
print(task10_result.collect())
[0]