【问题标题】:Consume large xml from hdfs with spark-streaming使用 spark-streaming 从 hdfs 使用大型 xml
【发布时间】:2018-04-23 08:35:36
【问题描述】:

我是listening 到hdfs 目录的xml 记录与spark-streaming-textFileStream()。问题是我的记录很大(而且是单行的);它们的大小可以接近 1G。

我愿意:

val xmlStream = ssc.textFileStream(monitoredDirectory).map { ("",_) }

但 spark 将我的文件拆分为更好的并行处理。 Xml 是一种不可分割的格式,我对文件的处理并没有很好地结束。

如何告诉 spark 不要拆分我的文件?还是有其他方法来处理大型 xml 文件?

【问题讨论】:

  • 可以添加当前使用的代码吗?
  • 代码,就这么简单。
  • 您使用的是什么版本的 spark?您是否试图简单地将文件作为整个文本文件读取?是否有理由使用火花流而不是结构化流?
  • 我使用的是 spark 1.6,所以没有结构化流。
  • 您是否已经尝试过spark-xml 库?

标签: xml scala hdfs spark-streaming


【解决方案1】:

在我看来,要管理大文件,流式传输并不是最佳解决方案。简单的方法是通过简单的方式管理它们

sc.textfile("newfileinthefolder", partition=1)

并使用文件夹中的侦听器调用此作业,但这样您会失去(或延迟)解决方案的实时计算功能。考虑一下您是否不需要近乎实时的功能。

另一个解决方案,但我对此不太自信,可以管理 StreamingContext 的 batchDuration。 在这种情况下,请注意您的流媒体生成的沿袭。 最后,看看this,databricks 资源是最合适的解决方案,有时

【讨论】:

  • 我明白,但这个特定的用例是关于流媒体的。我不明白你回答的最后一部分。
  • 当你定义 ssc 时,你必须设置 batchDuration(流数据将被分成批次的时间间隔)。也许 DStream 生成带有切割线的 RDD,因为 batchDuration 不够。如果您使用流式传输,请记住使用检查点,因为 batchDuration 越大,内存中需要更多空间。我可以给出的最后一个建议是关于 databrick 的 library 来管理 xml withspark。
  • 好的。我以为您在谈论有关批处理持续时间管理的更高阶概念,而不仅仅是批处理持续时间设置。我明白。好建议。
【解决方案2】:

用spark-xml,像建议的gtosto,像这样:

import org.apache.hadoop.io.{LongWritable, Text}
import com.databricks.spark.xml.XmlInputFormat

val conf = sc.hadoopConfiguration
conf.set(XmlInputFormat.START_TAG_KEY, "<xxx>")
conf.set(XmlInputFormat.END_TAG_KEY, "</xxx>")
org.apache.hadoop.fs.FileSystem.get(conf)

val xml = ssc.fileStream[LongWritable,Text,XmlInputFormat](monitoredDirectory,true,false)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-12-18
    • 2018-12-15
    • 2018-08-05
    • 1970-01-01
    • 2016-06-23
    • 1970-01-01
    • 2015-12-03
    • 2015-12-10
    相关资源
    最近更新 更多