【问题标题】:How to apply filters while reading all the json files from a folder using scala?如何在使用scala从文件夹中读取所有json文件时应用过滤器?
【发布时间】:2019-11-13 15:46:21
【问题描述】:

我有一个包含多个 json 文件的文件夹(first.json,second.json)。使用 scala,我将所有 jsonfiles 数据加载到 spark 的 rdd/dataset,然后对数据应用过滤器。

这里的问题是,如果我们有 600 个数据,那么我们需要将它们全部加载到 rdd/dataset 中,然后我们正在应用过滤器

寻找一种解决方案,我可以在从文件夹本身读取而不加载到 spark 内存中的同时过滤记录。

过滤是基于blockheight属性完成的。

每个文件中的Json结构:

first.json:

[{"IsFee":false,"BlockDateTime":"2015-10-14T09:02:46","Address":"0xe8fdc802e721426e0422d18d371ab59a41ddaeac","BlockHeight":381859,"Type":"IN","值":0.61609232637203584,"TransactionHash":"0xe6fc01ff633b4170e0c8f2df7db717e0608f8aaf62e6fbf65232a7009b53da4e","UserName":null,"ProjectName":null,"CreatedUser":null,"Id":0:,"CreatedUserId":0,"CreatedTime 08-26T22:32:45.2686137+05:30","UpdatedUserId":0,"UpdatedTime":"2019-08-26T22:32:45.2696126+05:30"},{"IsFee":false,"BlockDateTime" :"2015-10-14T09:02:46","Address":"0x52bc44d5378309ee2abf1539bf71de1b7d7be3b5","BlockHeight":381859,"Type":"OUT","Value":-0.61609232637203584,"TransactionHash":"0xe6fc01ff633b4170e0c8f2df7db717e0608f8aaf62e6fbf65232a7009b53da4e", "UserName":null,"ProjectName":null,"CreatedUser":null,"Id":0,"CreatedUserId":0,"CreatedTime":"2019-08-26T22:32:45.3141203+05:30", "UpdatedUserId":0,"UpdatedTime":"2019-08-26T22:32:45.3141203+05:30"}]

import org.apache.spark.SparkContext
import org.apache.spark.SparkContext._
import org.apache.spark.SparkConf
import org.apache.spark.sql.functions._
import org.apache.spark.sql._
import org.apache.spark.sql.types._

object BalanceAndTransactionDownload {
    def main(args: Array[String]) {
    val spark = SparkSession.builder.appName("xxx").getOrCreate()

    val currencyDataSchema = StructType(Array(
        StructField("Type", StringType, true),
        StructField("TransactionHash", StringType, true),
        StructField("BlockHeight", LongType, true),
        StructField("BlockDateTime", TimestampType, true),
        StructField("Value", DecimalType(38, 18), true),
        StructField("Address", StringType, true),
        StructField("IsFee", BooleanType, true)
    ))

    val projectAddressesFile = args(0)
    val blockJSONFilesContainer = args(1)
    val balanceFolderName = args(2)
    val downloadFolderName = args(3)
    val blockHeight = args(4)
    val projectAddresses = spark.read.option("multiline", "true").json(projectAddressesFile)
    val currencyDataFile = spark.read.option("multiline", "true").schema(currencyDataSchema).json(blockJSONFilesContainer) // This is where i want to filter out the data
    val filteredcurrencyData = currencyDataFile.filter(currencyDataFile("BlockHeight") <= blockHeight)
    filteredcurrencyData.join(projectAddresses, filteredcurrencyData("Address") === projectAddresses("address")).groupBy(projectAddresses("address")).agg(sum("Value").alias("Value")).repartition(1).write.option("header", "true").format("com.databricks.spark.csv").csv(balanceFolderName)
    filteredcurrencyData.join(projectAddresses, filteredcurrencyData("Address") === projectAddresses("address")).drop(projectAddresses("address")).drop(projectAddresses("CurrencyId")).drop(projectAddresses("Id")).repartition(1).write.option("header", "true").format("com.databricks.spark.csv").csv(downloadFolderName)
    }
}

【问题讨论】:

    标签: apache-spark-2.0


    【解决方案1】:

    应该对数据存储上的文件进行分区。你似乎是按块高度过滤的。所以你可以有多个文件夹,如:blockheight=1blockheight=2 等,并在这些文件夹中有 json 文件。在这种情况下,spark 不会读取所有 json 文件,而是会扫描所需的文件夹。

    【讨论】:

    • 感谢 Ravi 的提醒。但是我应该将文件夹重命名为 blockheight=1 还是简单地 1 , 2 .... 1000 并且在每个文件夹中我都可以拥有 json 文件,例如文件夹 1 将拥有 1.json ,文件夹 2 将拥有 2.json 。请让我知道这是否可行,或者我需要使用 blockheight=1 等等
    • blockheight=1 。只有当您以这种方式创建文件夹时,它们才会被谓词查询:blockheight。每个文件夹都可以有一个任意名称的 json; UUID 也有效。名称(1.json 或仅 1 或 abc.json 或仅 pqr)应该无关紧要。但请确保 json 文件与相应的文件夹 blockheight=&lt;&gt; 相关
    猜你喜欢
    • 1970-01-01
    • 2021-01-08
    • 1970-01-01
    • 2021-04-29
    • 1970-01-01
    • 2010-12-23
    • 2021-05-28
    相关资源
    最近更新 更多