Skip to content

Commit d2099d8

Browse files
Ken Takagiwagiwa
authored andcommitted
sorted the import following Spark coding convention
1 parent 5bac7ec commit d2099d8

File tree

1 file changed

+13
-32
lines changed

1 file changed

+13
-32
lines changed

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

Lines changed: 13 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -19,42 +19,28 @@ package org.apache.spark.streaming.api.python
1919

2020
import java.util.{List => JList, ArrayList => JArrayList, Map => JMap, Collections}
2121

22-
import org.apache.spark.api.java.{JavaSparkContext, JavaPairRDD, JavaRDD}
23-
import org.apache.spark.broadcast.Broadcast
22+
import scala.reflect.ClassTag
23+
2424
import org.apache.spark._
25-
import org.apache.spark.util.Utils
26-
import java.io._
27-
import scala.Some
28-
import org.apache.spark.streaming.Duration
29-
import scala.util.control.Breaks._
30-
import org.apache.spark.broadcast.Broadcast
31-
import scala.Some
32-
import org.apache.spark.streaming.Duration
3325
import org.apache.spark.rdd.RDD
34-
import org.apache.spark.api.python.PythonRDD
35-
36-
26+
import org.apache.spark.api.python._
27+
import org.apache.spark.broadcast.Broadcast
3728
import org.apache.spark.streaming.{Duration, Time}
3829
import org.apache.spark.streaming.dstream._
3930
import org.apache.spark.streaming.api.java._
40-
import org.apache.spark.rdd.RDD
41-
import org.apache.spark.api.python._
42-
import org.apache.spark.api.python.PairwiseRDD
43-
4431

45-
import scala.reflect.ClassTag
4632

4733

4834
class PythonDStream[T: ClassTag](
49-
parent: DStream[T],
50-
command: Array[Byte],
51-
envVars: JMap[String, String],
52-
pythonIncludes: JList[String],
53-
preservePartitoning: Boolean,
54-
pythonExec: String,
55-
broadcastVars: JList[Broadcast[Array[Byte]]],
56-
accumulator: Accumulator[JList[Array[Byte]]]
57-
) extends DStream[Array[Byte]](parent.ssc) {
35+
parent: DStream[T],
36+
command: Array[Byte],
37+
envVars: JMap[String, String],
38+
pythonIncludes: JList[String],
39+
preservePartitoning: Boolean,
40+
pythonExec: String,
41+
broadcastVars: JList[Broadcast[Array[Byte]]],
42+
accumulator: Accumulator[JList[Array[Byte]]])
43+
extends DStream[Array[Byte]](parent.ssc) {
5844

5945
override def dependencies = List(parent)
6046

@@ -146,8 +132,3 @@ DStream[(Long, Array[Byte])](prev.ssc){
146132
}
147133
val asJavaPairDStream : JavaPairDStream[Long, Array[Byte]] = JavaPairDStream.fromJavaDStream(this)
148134
}
149-
150-
151-
152-
153-

0 commit comments

Comments
 (0)