Skip to content

Commit 58591d2

Browse files
Ken Takagiwagiwa
Ken Takagiwa
authored andcommitted
reduceByKey is working
1 parent 0df7111 commit 58591d2

File tree

2 files changed

+12
-7
lines changed

2 files changed

+12
-7
lines changed
1.53 KB
Binary file not shown.

streaming/src/main/scala/org/apache/spark/streaming/api/python/PythonTransformedDStream.scala

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
/*
2+
13
package org.apache.spark.streaming.api.python
24
35
import org.apache.spark.Accumulator
@@ -10,11 +12,8 @@ import org.apache.spark.streaming.dstream.DStream
1012
1113
import scala.reflect.ClassTag
1214
13-
/**
14-
* Created by ken on 7/15/14.
15-
*/
1615
class PythonTransformedDStream[T: ClassTag](
17-
parents: Seq[DStream[T]],
16+
parent: DStream[T],
1817
command: Array[Byte],
1918
envVars: JMap[String, String],
2019
pythonIncludes: JList[String],
@@ -30,8 +29,14 @@ class PythonTransformedDStream[T: ClassTag](
3029
3130
//pythonDStream compute
3231
override def compute(validTime: Time): Option[RDD[Array[Byte]]] = {
33-
val parentRDDs = parents.map(_.getOrCompute(validTime).orNull).toSeq
34-
Some()
32+
33+
// val parentRDDs = parents.map(_.getOrCompute(validTime).orNull).toSeq
34+
// parents.map(_.getOrCompute(validTime).orNull).to
35+
// parent = parents.head.asInstanceOf[RDD]
36+
// Some()
3537
}
36-
val asJavaDStream = JavaDStream.fromDStream(this)
38+
39+
val asJavaDStream = JavaDStream.fromDStream(this)
3740
}
41+
42+
*/

0 commit comments

Comments
 (0)