【问题标题】:Apache Spark write to multiple outputs [different parquet schemas] without cachingApache Spark 写入多个输出 [不同的镶木地板模式],无需缓存
【发布时间】:2019-03-01 19:51:59
【问题描述】:

我想转换我的输入数据(XML 文件)并产生 3 个不同的输出。

每个输出都将采用 parquet 格式,并且具有不同的架构/列数。

目前在我的解决方案中,数据存储在RDD[Row] 中,其中每一行属于三种类型之一,并且具有不同数量的字段。我现在正在做的是缓存 RDD,然后过滤它(使用告诉我记录类型的字段)并使用以下方法保存数据:

var resultDF_1 = sqlContext.createDataFrame(filtered_data_1, schema_1)
resultDF_1.write.parquet(output_path_1)
...
// the same for filtered_data_2 and filtered_data_3

有什么办法可以做得更好,例如不要在内存中缓存整个数据?

在 MapReduce 中,我们有 MultipleOutputs 类,我们可以这样做:

MultipleOutputs.addNamedOutput(job, "data_type_1", DataType1OutputFormat.class, Void.class, Group.class);
MultipleOutputs.addNamedOutput(job, "data_type_2", DataType2OutputFormat.class, Void.class, Group.class);
MultipleOutputs.addNamedOutput(job, "data_type_3", DataType3OutputFormat.class, Void.class, Group.class);
...
MultipleOutputs<Void, Group> mos = new MultipleOutputs<>(context);
mos.write("data_type_1", null, myRecordGroup1, filePath1);
mos.write("data_type_2", null, myRecordGroup2, filePath2);
...

【问题讨论】:

  • 不确定我是否理解问题...您不想加载整个数据集?您是坚持使用 RDD 还是可以使用数据帧(性能应该更好 - 如果这是您的问题)?
  • 是的,我想加载整个数据集。我目前的问题是数据缓存。我有很多输入数据要处理。在我当前的解决方案中,我必须将一个 RDD 拆分为 3 个具有不同架构的独立 RDD。我使用过滤器功能,所以一开始我缓存了整个 RDD。问题的大致轮廓是:如何处理数据并保存到不同schem的不同parquet文件中。示例:1 条输入 XML 消息由请求和响应部分组成。您想将其转换为 parquet 并生成两个输出数据集:一个用于 RQ,另一个用于 RS

标签: apache-spark apache-spark-sql


【解决方案1】:

AFAIK,没有办法将一个 RDD 本身拆分为多个 RDD。这就是 Spark 的 DAG 的工作方式:只有子 RDD 从父 RDD 中提取数据。

但是,我们可以让多个子 RDD 从同一个父 RDD 中读取。为了避免重新计算父 RDD,除了缓存它别无他法。 我假设您想避免缓存,因为您害怕内存不足。我们可以通过将 RDD 持久化到 MEMORY_AND_DISK 来避免内存不足 (OOM) 问题,这样大的 RDD 就会溢出到磁盘必要时。

让我们从您的原始数据开始:

val allDataRDD = sc.parallelize(Seq(Row(1,1,1),Row(2,2,2),Row(3,3,3)))

我们可以先将它持久化到内存中,但在内存不足的情况下让它溢出到磁盘:

allDataRDD.persist(StorageLevel.MEMORY_AND_DISK)

然后我们创建 3 个 RDD 输出:

filtered_data_1 = allDataRDD.filter(_.get(1)==1) // //
filtered_data_2 = allDataRDD.filter(_.get(2)==1) // use your own filter funcs here
filtered_data_3 = allDataRDD.filter(_.get(3)==1) // //

然后我们编写输出:

var resultDF_1 = sqlContext.createDataFrame(filtered_data_1, schema_1)
resultDF_1.write.parquet(output_path_1)
var resultDF_2 = sqlContext.createDataFrame(filtered_data_2, schema_2)
resultDF_2.write.parquet(output_path_2)
var resultDF_3 = sqlContext.createDataFrame(filtered_data_3, schema_3)
resultDF_3.write.parquet(output_path_3)

如果您真的想避免多次传递,可以使用自定义分区程序来解决此问题。您可以将数据重新分区为 3 个分区,每个分区都有自己的任务,因此也有自己的输出文件/部分。需要注意的是,并行性将大大减少到 3 个线程/任务,并且还存在在单个分区中存储超过 2GB 数据的风险(Spark 每个分区有 2GB 的限制)。我不提供此方法的详细代码,因为我认为它不能编写具有不同架构的 parquet 文件。

【讨论】:

  • 感谢@LeightonRitchie 的反馈。是的,我想避免由于不必要的资源浪费(内存和可能的磁盘)而进行缓存。我们目前正在移动我们的旧项目,该项目是作为 MapReduce 作业编写的。在 MR 方法中,我们使用MultipleOutputs,并且数据不会缓存在任何地方。我们正在寻找如何在 Spark 中实现类似的功能。
  • 当您使用 RDD 级别抽象时,Spark 支持所有 Hadoop 输出格式,但在数据集级别,写入逻辑是抽象的。也许使用 RDD 会起作用。它也需要编写以拼花格式编写的逻辑。
猜你喜欢
  • 2017-03-17
  • 1970-01-01
  • 2022-12-01
  • 2021-03-15
  • 2019-01-08
  • 2019-06-02
  • 2015-11-20
  • 1970-01-01
  • 2016-08-12
相关资源
最近更新 更多