File tree Expand file tree Collapse file tree 4 files changed +2
-10
lines changed
main/scala/org/apache/spark/streaming/kafka
test/java/org/apache/spark/streaming/kafka Expand file tree Collapse file tree 4 files changed +2
-10
lines changed Original file line number Diff line number Diff line change @@ -196,8 +196,8 @@ private class KafkaTestUtils extends Logging {
196196 props
197197 }
198198
199- // A simplified version of scalatest eventually, rewritten here is to avoid adding extra test
200- // dependency
199+ // A simplified version of scalatest eventually, rewritten here to avoid adding extra test
200+ // dependency
201201 def eventually [T ](timeout : Time , interval : Time )(func : => T ): T = {
202202 def makeAttempt (): Either [Throwable , T ] = {
203203 try {
Original file line number Diff line number Diff line change @@ -47,7 +47,6 @@ public class JavaDirectKafkaStreamSuite implements Serializable {
4747 public void setUp () {
4848 kafkaTestUtils = new KafkaTestUtils ();
4949 kafkaTestUtils .setup ();
50- System .clearProperty ("spark.driver.port" );
5150 SparkConf sparkConf = new SparkConf ()
5251 .setMaster ("local[4]" ).setAppName (this .getClass ().getSimpleName ());
5352 ssc = new JavaStreamingContext (sparkConf , Durations .milliseconds (200 ));
@@ -60,8 +59,6 @@ public void tearDown() {
6059 ssc = null ;
6160 }
6261
63- System .clearProperty ("spark.driver.port" );
64-
6562 if (kafkaTestUtils != null ) {
6663 kafkaTestUtils .teardown ();
6764 kafkaTestUtils = null ;
Original file line number Diff line number Diff line change @@ -43,7 +43,6 @@ public class JavaKafkaRDDSuite implements Serializable {
4343 public void setUp () {
4444 kafkaTestUtils = new KafkaTestUtils ();
4545 kafkaTestUtils .setup ();
46- System .clearProperty ("spark.driver.port" );
4746 SparkConf sparkConf = new SparkConf ()
4847 .setMaster ("local[4]" ).setAppName (this .getClass ().getSimpleName ());
4948 sc = new JavaSparkContext (sparkConf );
@@ -55,7 +54,6 @@ public void tearDown() {
5554 sc .stop ();
5655 sc = null ;
5756 }
58- System .clearProperty ("spark.driver.port" );
5957
6058 if (kafkaTestUtils != null ) {
6159 kafkaTestUtils .teardown ();
Original file line number Diff line number Diff line change @@ -48,7 +48,6 @@ public class JavaKafkaStreamSuite implements Serializable {
4848 public void setUp () {
4949 kafkaTestUtils = new KafkaTestUtils ();
5050 kafkaTestUtils .setup ();
51- System .clearProperty ("spark.driver.port" );
5251 SparkConf sparkConf = new SparkConf ()
5352 .setMaster ("local[4]" ).setAppName (this .getClass ().getSimpleName ());
5453 ssc = new JavaStreamingContext (sparkConf , new Duration (500 ));
@@ -61,8 +60,6 @@ public void tearDown() {
6160 ssc = null ;
6261 }
6362
64- System .clearProperty ("spark.driver.port" );
65-
6663 if (kafkaTestUtils != null ) {
6764 kafkaTestUtils .teardown ();
6865 kafkaTestUtils = null ;
You can’t perform that action at this time.
0 commit comments