Skip to content

Commit f1996d3

Browse files
committed
[SPARK-2548][HOTFIX] Removed use of o.a.s.streaming.Durations
1 parent cdcf546 commit f1996d3

File tree

2 files changed

+4
-4
lines changed

2 files changed

+4
-4
lines changed

examples/src/main/java/org/apache/spark/examples/streaming/JavaNetworkWordCount.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@
2525
import org.apache.spark.api.java.function.Function2;
2626
import org.apache.spark.api.java.function.PairFunction;
2727
import org.apache.spark.api.java.StorageLevels;
28-
import org.apache.spark.streaming.Durations;
28+
import org.apache.spark.streaming.Duration;
2929
import org.apache.spark.streaming.api.java.JavaDStream;
3030
import org.apache.spark.streaming.api.java.JavaPairDStream;
3131
import org.apache.spark.streaming.api.java.JavaReceiverInputDStream;
@@ -57,7 +57,7 @@ public static void main(String[] args) {
5757

5858
// Create the context with a 1 second batch size
5959
SparkConf sparkConf = new SparkConf().setAppName("JavaNetworkWordCount");
60-
JavaStreamingContext ssc = new JavaStreamingContext(sparkConf, Durations.seconds(1));
60+
JavaStreamingContext ssc = new JavaStreamingContext(sparkConf, new Duration(1000));
6161

6262
// Create a JavaReceiverInputDStream on target ip:port and count the
6363
// words in input stream of \n delimited text (eg. generated by 'nc')

examples/src/main/java/org/apache/spark/examples/streaming/JavaRecoverableNetworkWordCount.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@
3232
import org.apache.spark.api.java.function.FlatMapFunction;
3333
import org.apache.spark.api.java.function.Function2;
3434
import org.apache.spark.api.java.function.PairFunction;
35-
import org.apache.spark.streaming.Durations;
35+
import org.apache.spark.streaming.Duration;
3636
import org.apache.spark.streaming.Time;
3737
import org.apache.spark.streaming.api.java.JavaDStream;
3838
import org.apache.spark.streaming.api.java.JavaPairDStream;
@@ -83,7 +83,7 @@ private static JavaStreamingContext createContext(String ip,
8383
}
8484
SparkConf sparkConf = new SparkConf().setAppName("JavaRecoverableNetworkWordCount");
8585
// Create the context with a 1 second batch size
86-
JavaStreamingContext ssc = new JavaStreamingContext(sparkConf, Durations.seconds(1));
86+
JavaStreamingContext ssc = new JavaStreamingContext(sparkConf, new Duration(1000));
8787
ssc.checkpoint(checkpointDirectory);
8888

8989
// Create a socket stream on target ip:port and count the

0 commit comments

Comments
 (0)