【问题标题】:How to read multiple directories in s3 in spark Scala?如何在 Spark Scala 中读取 s3 中的多个目录?
【发布时间】:2018-03-08 08:49:18
【问题描述】:

我在 s3 中有以下格式的目录,

 <base-directory>/users/users=20180303/hour=0/<parquet files>
 <base-directory>/users/users=20180303/hour=1/<parquet files>
 ....
 <base-directory>/users/users=20180302/hour=<0 to 23>/<parquet files>
 <base-directory>/users/users=20180301/hour=<0 to 23>/<parquet files>
 ....
 <base-directory>/users/users=20180228/hour=<0 to 23>/<parquet files>

基本上我在日常目录中有每小时的子目录。

现在我想处理过去 30 天的 parquet 文件。

我已经尝试过,

 val df = sqlContext.read.option("header", "true")
    .parquet(<base-directory> + File.separator + "users" + File.separator)
    .where(col("users").between(startDate, endDate))

endDate 和 startDate 相隔 30 天,格式为 yyyymmdd。

上述解决方案未提供正确的目录子集。我做错了什么?

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    where 函数用于dataframe 中的过滤行。您正在使用它从 s3 读取 parquet 文件。 所以整个概念是错误的

    相反,您可以在 startDate 和 endDate 之间创建一个路径数组,并将其传递给 sqlContext read api

    从编程上讲,您可以执行以下操作(它们只是伪代码)

    val listBuffer = new ListBuffer[String]
    for(date <- startDate to endDate)
      listBuffer.append(<base-directory> + File.separator + "users" + File.separator+"users="+date)
    
    val df = sqlContext.read.option("header", "true").parquet(listBuffer: _*)
    

    【讨论】:

    • 我很高兴@abhijeet 感谢您的支持和接受:)
    • ..是否可以通过电话与您交谈或通过 IM 聊天...我正在解决一个棘手的问题。告诉我。
    • 你可以在 stackoverflow @abhijeet 上提问,不仅是我,每个人都会帮助你
    • ..好的。我请求您查看以下链接,stackoverflow.com/questions/49237513/…
    • 我请你看看下面的链接stackoverflow.com/questions/49251699/…
    猜你喜欢
    • 2021-07-31
    • 2018-04-18
    • 2021-04-27
    • 1970-01-01
    • 2021-08-15
    • 2020-01-02
    • 2020-02-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多