【问题标题】:Spark Streaming from socket does not work with reduce operation来自套接字的 Spark Streaming 不适用于 reduce 操作
【发布时间】:2013-06-26 01:56:36
【问题描述】:

我正在尝试在我的本地机器上运行一个简单的 Spark-Streaming 示例。
我有一个将 As/Bs/Cs 写入套接字的线程:

serverSocket = new ServerSocket(Constants.PORT);
s1 = serverSocket.accept();
while(true) {
    Thread.sleep(random.nextInt(100));
    String character = alphabet.get(random.nextInt(alphabet.size())) ;
    PrintWriter out = new PrintWriter(s1.getOutputStream());
    out.println(character);
    out.flush();
}

我的主程序,我尝试计算 As/Bs/Cs 的数量如下所示(没有减少步骤):

public static void main(String[] args) {
    // start socket writer thread
    System.setProperty("spark.cleaner.ttl", "10000");
    JavaSparkContext sc = new JavaSparkContext(
            "local", 
            "Test",
            Constants.SPARK_HOME, 
            new String[]{"target/spark-standalone-0.0.1-SNAPSHOT.jar"});
    Duration batchDuration = new Duration(TIME_WINDOW_MS);
    JavaStreamingContext streamingContext = new JavaStreamingContext(sc, batchDuration);
    JavaDStream<String> stream = streamingContext.socketTextStream("localhost", Constants.PORT);
    stream.print();
    JavaPairDStream<String, Long> texts = stream.map(new PairFunction<String, String, Long>() {

            @Override
            public Tuple2<String, Long> call(String t) throws Exception {
                return new Tuple2<String, Long>("batchCount" + t, 1l);
            }

        });
     texts.print();
     streamingContext.checkpoint("checkPointDir");
     streamingContext.start();

在这种情况下,一切正常(批处理的示例输出):

Time: 1372413296000 ms
-------------------------------------------
B
A
B
C
C
C
A
B
C
C
...

-------------------------------------------
Time: 1372413296000 ms
-------------------------------------------
(batchCountB,1)
(batchCountA,1)
(batchCountB,1)
(batchCountC,1)
(batchCountC,1)
(batchCountC,1)
(batchCountA,1)
(batchCountB,1)
(batchCountC,1)
(batchCountC,1)
...

但是,如果我在地图之后添加缩减步骤,它就不再起作用了。此代码位于 texts.print() 之后

JavaPairDStream<String, Long> reduced = texts.reduceByKeyAndWindow(new Function2<Long, Long, Long>() {

    @Override
    public Long call(Long t1, Long t2) throws Exception {
        return t1 + t2;
    }
    }, new Duration(TIME_WINDOW_MS));
reduced.print();

在这种情况下,我只得到第一个“stream”变量和“texts”变量的输出,而 reduce 没有任何输出。在第一次批处理之后也没有任何反应。我还将 spark 日志级别设置为 DEBUG,但没有遇到任何异常或其他奇怪的事情。

这里发生了什么?为什么我会被锁定?

【问题讨论】:

    标签: java apache-spark streaming bigdata


    【解决方案1】:

    仅作记录:我在 Spark 用户组中得到了答案。
    错误是必须使用

    "local[2]"
    

    而不是

    "local"
    

    作为参数来实例化 Spark 上下文,以启用并发处理。

    【讨论】:

      猜你喜欢
      • 2011-01-23
      • 2017-07-10
      • 1970-01-01
      • 2016-11-03
      • 1970-01-01
      • 2014-04-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多