【问题标题】:Apache Spark read file as a stream from HDFSApache Spark 将文件作为流从 HDFS 读取
【发布时间】:2017-06-14 00:21:00
【问题描述】:

如何使用 Apache Spark Java 从 hdfs 以流的形式读取文件? 我不想读取整个文件,我想要文件流以便在满足某些条件时停止读取文件,我该如何使用 Apache Spark 来做到这一点?

【问题讨论】:

  • 这个例子与我的问题无关。
  • 你能更好地解释你想要达到的目标吗?为什么需要它作为流(而不是简单地将其作为 RDD/Dataframe 读取)?您是否在问如何让火花流读取 HDFS 目录的内容并在完成时停止(而不是等待下一个时间段)?您是在谈论 DStream 还是结构化流式传输?
  • 问题是,例如当您尝试部分读取 Parquet 文件时会发生什么?我想说的是:当文件系统支持在完全下载之前无法解码的专有文件格式时,(对于 Hadoop 开发人员)制作这样的功能是否有意义?但这只是一个想法。
  • Assaf Mendelson,我想逐字节读取文件并在满足某些条件后停止读取,例如找到了一些符号...有没有可能,或者 NameNode 总是会查找所有文件块?

标签: java apache-spark hdfs


【解决方案1】:

您可以使用 ssc 方法使用流式 HDFS 文件

val ssc = new StreamingContext(sparkConf, Seconds(batchTime))

val dStream = ssc.fileStream[LongWritable, Text, TextInputFormat]( streamDirectory, (x: Path) => true, newFilesOnly = false)

使用上面的api param filter 过滤路径的函数。

如果您的条件不是文件路径/名称并且基于数据,那么如果条件满足,您需要停止流式传输上下文。

为此,您需要使用线程实现, 1)在一个线程中,您需要继续检查流上下文是否已停止,如果 ssc 停止,则通知其他线程等待并创建新的流上下文。

2) 在第二个线程中,您需要检查条件,如果条件满足则停止流式传输上下文。

如果您需要解释,请告诉我。

【讨论】:

  • 我遇到的问题,例如两千个文件,我只想从每个文件中读取 N 行(从几行到数十亿行)。您的解决方案将非常昂贵。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-09-11
  • 1970-01-01
  • 2021-06-18
  • 1970-01-01
  • 2021-06-25
  • 1970-01-01
  • 2018-12-15
相关资源
最近更新 更多