【发布时间】:2018-01-28 09:38:09
【问题描述】:
有没有人知道怎么设置`
streamingContext.fileStream [KeyClass, ValueClass, InputFormatClass] (dataDirectory)
实际使用二进制文件。
在哪里可以找到所有 inputformatClass ?文档没有给出 链接。我想 ValueClass 与 inputformatClass 不知何故。
在使用binaryfiles方法的非流式版本中,我可以得到 每个文件的字节数组。有没有办法我可以得到相同的 火花流?如果没有,我在哪里可以找到这些详细信息。意思是 支持的输入格式和它产生的值类。终于可以了 随便选一个 KeyClass,这些元素不是都连接了吗?
如果有人澄清了该方法的使用。
EDIT1
我尝试了以下方法:
val bfiles = ssc.fileStreamBytesWritable, BytesWritable, SequenceFileAsBinaryInputFormat
但是编译器会这样抱怨:
[error] /xxxxxxxxx/src/main/scala/EstimatorStreamingApp.scala:14: type arguments [org.apache.hadoop.io.BytesWritable,org.apache.hadoop.io.BytesWritable,org.apache.hadoop.mapred.SequenceFileAsBinaryInputFormat] conform to the bounds of none of the overloaded alternatives of
[error] value fileStream: [K, V, F <: org.apache.hadoop.mapreduce.InputFormat[K,V]](directory: String, filter: org.apache.hadoop.fs.Path => Boolean, newFilesOnly: Boolean, conf: org.apache.hadoop.conf.Configuration)(implicit evidence$10: scala.reflect.ClassTag[K], implicit evidence$11: scala.reflect.ClassTag[V], implicit evidence$12: scala.reflect.ClassTag[F])org.apache.spark.streaming.dstream.InputDStream[(K, V)] <and> [K, V, F <: org.apache.hadoop.mapreduce.InputFormat[K,V]](directory: String, filter: org.apache.hadoop.fs.Path => Boolean, newFilesOnly: Boolean)(implicit evidence$7: scala.reflect.ClassTag[K], implicit evidence$8: scala.reflect.ClassTag[V], implicit evidence$9: scala.reflect.ClassTag[F])org.apache.spark.streaming.dstream.InputDStream[(K, V)] <and> [K, V, F <: org.apache.hadoop.mapreduce.InputFormat[K,V]](directory: String)(implicit evidence$4: scala.reflect.ClassTag[K], implicit evidence$5: scala.reflect.ClassTag[V], implicit evidence$6: scala.reflect.ClassTag[F])org.apache.spark.streaming.dstream.InputDStream[(K, V)]
[error] val bfiles = ssc.fileStream[BytesWritable, BytesWritable, SequenceFileAsBinaryInputFormat]("/xxxxxxxxx/Casalini_streamed")
我做错了什么?
【问题讨论】:
-
你检查here所有hadoop输入格式了吗?如果您的数据集是序列文件二进制/原始格式,请尝试 SequenceFileAsBinaryInputFormat
-
@squid 我已经更新了帖子,我根据您的输入添加了代码。请看一下,我仍然有一些编译器问题。虽然我认为我写的是正确的。
标签: apache-spark spark-streaming