【发布时间】:2018-05-25 06:25:27
【问题描述】:
我正在查看这个火花样本:
https://spark.apache.org/docs/latest/streaming-programming-guide.html#a-quick-example
我想实现一个“echo”应用程序,一个火花流软件,它从一个套接字读取一些字符并在同一个套接字上输出相同的字符(这是我真正问题的简化,当然我有一些处理在输入上做,但我们现在不关心这个)。
我尝试按照本指南实施 CustomReceiver: https://spark.apache.org/docs/2.3.0/streaming-custom-receivers.html
我添加了一个方法getSocket:
def getSocket() : Socket = {
socket
}
我试着像这样调用它:
val receiver = new SocketReceiver2("localhost", 9999, StorageLevel.MEMORY_AND_DISK_2)
val lines = ssc.receiverStream(receiver)
lines.foreachRDD {
rdd => rdd.foreachPartition { partitionOfRecords =>
val os = receiver.getSocket().getOutputStream();
partitionOfRecords.foreach(record => os.write(record.getBytes()))
}
}
但我收到 Task not serializable 错误。 (正如 T.Gaweda 指出的那样,这是意料之中的)。所以下一步就是开始使用累加器...
有没有更简单的方法可以在 Spark Streaming 中执行我的“回声”应用程序?
(我真的需要使用 Kafka(hdfs、hive...)从一个简单的 java 应用程序来回发送数据吗?)
【问题讨论】:
标签: apache-spark spark-streaming