Skip to content

Commit 3888fe3

Browse files
committed
Remove setProperty use in LocalJavaStreamingContext
1 parent 4f4031d commit 3888fe3

File tree

5 files changed

+30
-10
lines changed

5 files changed

+30
-10
lines changed

external/flume/src/test/java/org/apache/spark/streaming/LocalJavaStreamingContext.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.spark.streaming;
1919

20+
import org.apache.spark.SparkConf;
2021
import org.apache.spark.streaming.api.java.JavaStreamingContext;
2122
import org.junit.After;
2223
import org.junit.Before;
@@ -27,8 +28,11 @@ public abstract class LocalJavaStreamingContext {
2728

2829
@Before
2930
public void setUp() {
30-
System.setProperty("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
31-
ssc = new JavaStreamingContext("local[2]", "test", new Duration(1000));
31+
SparkConf conf = new SparkConf()
32+
.setMaster("local[2]")
33+
.setAppName("test")
34+
.set("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
35+
ssc = new JavaStreamingContext(conf, new Duration(1000));
3236
ssc.checkpoint("checkpoint");
3337
}
3438

external/mqtt/src/test/java/org/apache/spark/streaming/LocalJavaStreamingContext.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.spark.streaming;
1919

20+
import org.apache.spark.SparkConf;
2021
import org.apache.spark.streaming.api.java.JavaStreamingContext;
2122
import org.junit.After;
2223
import org.junit.Before;
@@ -27,8 +28,11 @@ public abstract class LocalJavaStreamingContext {
2728

2829
@Before
2930
public void setUp() {
30-
System.setProperty("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
31-
ssc = new JavaStreamingContext("local[2]", "test", new Duration(1000));
31+
SparkConf conf = new SparkConf()
32+
.setMaster("local[2]")
33+
.setAppName("test")
34+
.set("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
35+
ssc = new JavaStreamingContext(conf, new Duration(1000));
3236
ssc.checkpoint("checkpoint");
3337
}
3438

external/twitter/src/test/java/org/apache/spark/streaming/LocalJavaStreamingContext.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.spark.streaming;
1919

20+
import org.apache.spark.SparkConf;
2021
import org.apache.spark.streaming.api.java.JavaStreamingContext;
2122
import org.junit.After;
2223
import org.junit.Before;
@@ -27,8 +28,11 @@ public abstract class LocalJavaStreamingContext {
2728

2829
@Before
2930
public void setUp() {
30-
System.setProperty("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
31-
ssc = new JavaStreamingContext("local[2]", "test", new Duration(1000));
31+
SparkConf conf = new SparkConf()
32+
.setMaster("local[2]")
33+
.setAppName("test")
34+
.set("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
35+
ssc = new JavaStreamingContext(conf, new Duration(1000));
3236
ssc.checkpoint("checkpoint");
3337
}
3438

external/zeromq/src/test/java/org/apache/spark/streaming/LocalJavaStreamingContext.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.spark.streaming;
1919

20+
import org.apache.spark.SparkConf;
2021
import org.apache.spark.streaming.api.java.JavaStreamingContext;
2122
import org.junit.After;
2223
import org.junit.Before;
@@ -27,8 +28,11 @@ public abstract class LocalJavaStreamingContext {
2728

2829
@Before
2930
public void setUp() {
30-
System.setProperty("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
31-
ssc = new JavaStreamingContext("local[2]", "test", new Duration(1000));
31+
SparkConf conf = new SparkConf()
32+
.setMaster("local[2]")
33+
.setAppName("test")
34+
.set("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
35+
ssc = new JavaStreamingContext(conf, new Duration(1000));
3236
ssc.checkpoint("checkpoint");
3337
}
3438

streaming/src/test/java/org/apache/spark/streaming/LocalJavaStreamingContext.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.spark.streaming;
1919

20+
import org.apache.spark.SparkConf;
2021
import org.apache.spark.streaming.api.java.JavaStreamingContext;
2122
import org.junit.After;
2223
import org.junit.Before;
@@ -27,8 +28,11 @@ public abstract class LocalJavaStreamingContext {
2728

2829
@Before
2930
public void setUp() {
30-
System.setProperty("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
31-
ssc = new JavaStreamingContext("local[2]", "test", new Duration(1000));
31+
SparkConf conf = new SparkConf()
32+
.setMaster("local[2]")
33+
.setAppName("test")
34+
.set("spark.streaming.clock", "org.apache.spark.streaming.util.ManualClock");
35+
ssc = new JavaStreamingContext(conf, new Duration(1000));
3236
ssc.checkpoint("checkpoint");
3337
}
3438

0 commit comments

Comments
 (0)