Skip to content

Commit

Permalink
WIP: solved partitioned and None is not recognized
Browse files Browse the repository at this point in the history
  • Loading branch information
giwa committed Sep 20, 2014
1 parent f67cf57 commit 1e126bf
Show file tree
Hide file tree
Showing 2 changed files with 17 additions and 0 deletions.
16 changes: 16 additions & 0 deletions python/pyspark/streaming/dstream.py
Original file line number Diff line number Diff line change
Expand Up @@ -245,6 +245,8 @@ def takeAndPrint(rdd, time):
taken = rdd.take(11)
print "-------------------------------------------"
print "Time: %s" % (str(time))
print rdd.glom().collect()
print "-------------------------------------------"
print "-------------------------------------------"
for record in taken[:10]:
print record
Expand Down Expand Up @@ -426,6 +428,20 @@ def saveAsTextFile(rdd, time):
# TODO: implemtnt rightOuterJoin


# TODO: implement groupByKey
# TODO: impelment union
# TODO: implement cache
# TODO: implement persist
# TODO: implement repertitions
# TODO: implement saveAsTextFile
# TODO: implement cogroup
# TODO: implement join
# TODO: implement countByValue
# TODO: implement leftOuterJoin
# TODO: implemtnt rightOuterJoin



class PipelinedDStream(DStream):
def __init__(self, prev, func, preservesPartitioning=False):
if not isinstance(prev, PipelinedDStream) or not prev._is_pipelinable():
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -205,3 +205,4 @@ class PythonTransformedDStream(
//val asJavaPairDStream : JavaPairDStream[Long, Array[Byte]] = JavaPairDStream.fromJavaDStream(this)
}
*/

0 comments on commit 1e126bf

Please sign in to comment.