【问题标题】:Reading binaryFile with Spark Streaming使用 Spark Streaming 读取 binaryFile
【发布时间】: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


【解决方案1】:

我终于可以编译了。

编译问题出在导入中。我用过

  • 导入 org.apache.hadoop.mapred.SequenceFileAsBinaryInputFormat

我换成

  • 导入 org.apache.hadoop.mapreduce.lib.input.SequenceFileAsBinaryInputFormat

然后就可以了。但是我不知道为什么。我不明白两个层次结构之间的区别。这两个文件似乎具有相同的内容。所以很难说。如果有人可以在这里帮助澄清这一点,我认为这将有很大帮助

【讨论】:

    【解决方案2】:

    点击链接阅读所有hadoop input formats

    我找到了here 有据可查的关于序列文件格式的答案。

    由于导入不匹配,您正面临编译问题。 Hadoop Mapred vs mapreduce

    例如

    Java

    JavaPairInputDStream<Text,BytesWritable> dstream=
            sc.fileStream("/somepath",org.apache.hadoop.io.Text.class,
            org.apache.hadoop.io.BytesWritable.class,
        org.apache.hadoop.mapreduce.lib.input.SequenceFileAsBinaryInputFormat.class);
    

    我没有在 scala 中尝试过,但应该是类似的;

    val dstream = sc.fileStream("/somepath", 
            classOf[org.apache.hadoop.io.Text], classOf[org.apache.hadoop.io.BytesWritable],
            classOf[org.apache.hadoop.mapreduce.lib.input.SequenceFileAsBinaryInputFormat] ) ;
    

    【讨论】:

      猜你喜欢
      • 2016-09-26
      • 2015-12-03
      • 2015-04-27
      • 1970-01-01
      • 1970-01-01
      • 2019-06-24
      • 2018-01-17
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多