【发布时间】: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