【发布时间】:2016-05-09 12:39:16
【问题描述】:
我正在开发一个 Spark Streaming 应用程序,我需要在其中使用来自 Python 中的两个服务器的输入流,每个服务器每秒向 Spark 上下文发送一条 JSON 消息。
我的问题是,如果我只对一个流执行操作,那么一切正常。但是,如果我有来自不同服务器的两个流,那么 Spark 在它可以打印任何东西之前就冻结了,并且只有在两个服务器都发送了它们必须发送的所有 JSON 消息时才重新开始工作(当它检测到 'socketTextStream 没有接收到数据。
这是我的代码:
JavaReceiverInputDStream<String> streamData1 = ssc.socketTextStream("localhost",996,
StorageLevels.MEMORY_AND_DISK_SER);
JavaReceiverInputDStream<String> streamData2 = ssc.socketTextStream("localhost", 9995,StorageLevels.MEMORY_AND_DISK_SER);
JavaPairDStream<Integer, String> dataStream1= streamData1.mapToPair(new PairFunction<String, Integer, String>() {
public Tuple2<Integer, String> call(String stream) throws Exception {
Tuple2<Integer,String> streamPair= new Tuple2<Integer, String>(1, stream);
return streamPair;
}
});
JavaPairDStream<Integer, String> dataStream2= streamData2.mapToPair(new PairFunction<String, Integer, String>() {
public Tuple2<Integer, String> call(String stream) throws Exception {
Tuple2<Integer,String> streamPair= new Tuple2<Integer, String>(2, stream);
return streamPair;
}
});
dataStream2.print(); //for example
请注意,没有 ERROR 消息,Spark 在启动上下文后简单冻结,当我从端口获取 JSON 消息时,它没有显示任何内容。
非常感谢。
【问题讨论】:
标签: java sockets apache-spark spark-streaming