Spark 1.6 Java版DataFrame转JSON数组性能调优及序列化问题
解决Spark 1.6+Java中DataFrame转JSON数组的性能与序列化问题
让我来帮你搞定这两个头疼的问题:一是Task not serializable错误,二是大表下reduce拼接JSON的性能瓶颈。
问题根源拆解
- 序列化错误:你在
commonTableDF.toJavaRDD().foreach()里直接引用了hiveContext、sc、dfs这些对象——它们都是Driver端的上下文实例,没法被序列化传到Executor节点执行,自然就抛出序列化异常了。 - reduce性能灾难:直接用
reduce把所有JSON字符串拼在一起,本质是把全量数据拉到单个节点做字符串拼接,每次拼接都会产生新的字符串对象(Java字符串不可变),时间复杂度是O(n²),588KB的数据都要几十分钟,大表肯定扛不住。
针对性解决方案
1. 修复序列化错误:把循环移到Driver端
表名列表是小数据集,完全没必要放到RDD里分布式遍历。直接在Driver端用普通Java循环处理每个表名,这样就不会在闭包里引用那些不可序列化的上下文对象了。
2. 优化JSON拼接:分区局部处理+Driver端合并
放弃全局reduce的思路,改用“分而治之”的策略:
- 先让每个Executor节点处理自己分区内的JSON字符串,把分区里的元素用逗号拼接成一个片段
- 再把所有分区的片段收集到Driver端,把这些片段用逗号连接,最后包裹
[]形成完整的JSON数组 - 这种方式把拼接压力分散到各个节点,避免单节点瓶颈,同时减少字符串复制的开销
另外,别用count()判断DataFrame是否为空——count()会触发一次全表扫描,改用jsonRDD.isEmpty()更高效,而且空RDD也不会生成有效输出文件。
完整修正代码
import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.compress.GzipCodec; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.sql.DataFrame; import org.apache.spark.sql.hive.HiveContext; import java.util.ArrayList; import java.util.Arrays; import java.util.List; public class JsonArrayConverter { public void convertDataFramesToJsonArrays(String[] cmnTableNames, HiveContext hiveContext, String snapshotId, String hdfsPath, String rptName, org.apache.spark.SparkContext sc, org.apache.hadoop.fs.FileSystem dfs) throws Exception { if (cmnTableNames != null && cmnTableNames.length > 0) { List<String> commonTableList = Arrays.asList(cmnTableNames); // 直接在Driver端遍历表名,彻底避开Executor端的序列化问题 for (String cmnTableName : commonTableList) { DataFrame cmnTableContent; // 构建查询SQL if (cmnTableName.contains("PTR_security_t")) { cmnTableContent = hiveContext.sql("SELECT * FROM " + cmnTableName + " where fbrn04_snapshot_d = '" + snapshotId + "'"); } else { cmnTableContent = hiveContext.sql("SELECT * FROM " + cmnTableName); } String cmnTable = cmnTableName.substring(cmnTableName.lastIndexOf(".") + 1); String cmnStgTblDir = hdfsPath + "/staging/" + rptName + "/common/" + cmnTable; String cmnStgMrgdDir = cmnStgTblDir + "/mergedfile"; // 先清理已存在的目标目录 Path mergedPath = new Path(cmnStgMrgdDir); if (dfs.exists(mergedPath)) { dfs.delete(mergedPath, true); } JavaRDD<String> jsonRDD = cmnTableContent.toJSON().toJavaRDD(); // 仅当RDD非空时处理 if (!jsonRDD.isEmpty()) { // 每个分区内拼接JSON字符串,生成局部片段 JavaRDD<String> partitionedJsonFragments = jsonRDD.mapPartitions(iter -> { StringBuilder sb = new StringBuilder(); boolean isFirstElement = true; while (iter.hasNext()) { if (!isFirstElement) { sb.append(","); } sb.append(iter.next()); isFirstElement = false; } return Arrays.asList(sb.toString()).iterator(); }); // 收集所有分区片段到Driver端,合并成完整JSON数组 List<String> allFragments = partitionedJsonFragments.collect(); StringBuilder finalJsonArray = new StringBuilder("["); boolean isFirstFragment = true; for (String fragment : allFragments) { if (!isFirstFragment) { finalJsonArray.append(","); } finalJsonArray.append(fragment); isFirstFragment = false; } finalJsonArray.append("]"); // 生成最终RDD并保存为压缩文件 List<String> outputList = new ArrayList<>(); outputList.add(finalJsonArray.toString()); JavaRDD<String> finalOutputRDD = sc.parallelize(outputList); finalOutputRDD.coalesce(1).saveAsTextFile(cmnStgMrgdDir, GzipCodec.class); // 设置文件权限(补充你之前截断的重命名逻辑) Path partFile = new Path(cmnStgMrgdDir + "/part-00000.gz"); if (dfs.exists(partFile)) { org.apache.hadoop.fs.FileStatus fileStatus = dfs.getFileStatus(partFile); dfs.setPermission(fileStatus.getPath(), org.apache.hadoop.fs.FsPermission.createImmutable((short) 0770)); // 这里可以添加重命名逻辑,比如改成更友好的文件名 // dfs.rename(partFile, new Path(cmnStgTblDir + "/"+ cmnTable +".json.gz")); } } } } } }
额外优化小贴士
- 避免SQL注入风险:别直接把变量拼进SQL字符串,建议用参数化查询(比如Spark的
hiveContext.sql支持?占位符),既安全又好维护。 - 减少不必要的Action:原代码里的
count()会触发一次全表扫描,改用jsonRDD.isEmpty()可以节省大量计算资源。 - 调整Spark资源配置:如果处理超大表,适当调大Executor的内存(
--executor-memory)和核心数(--executor-cores),提升并行处理能力。
内容的提问来源于stack exchange,提问作者Irthiza Khan
相关产品推荐
相关产品推荐

