【发布时间】:2020-08-22 06:25:00
【问题描述】:
这是一个不同的问题,我正在尝试根据列数过滤 RDD 中的记录。这更像是文件处理。
我在 Pyspark 中也写过相同的内容,并且我看到记录正在正确过滤。 当我在 Java 中尝试时,有效记录将进入错误文件夹。
下载错误文件并使用 AWK 进行验证,发现它们有 996 列,但仍然在错误中被过滤掉。
在 python 中,过滤的文件的确切数量是错误文件。
下面是片段。
JavaRDD<String> inputDataRDD = sc.textFile(args[0]+"/"+args[1], 5000);
int columnLength = Integer.parseInt(args[3]);
inputDataRDD
.filter(filterData -> filterData.split("\t").length == columnLength)
.coalesce(1)
.saveAsTextFile(args[2]+"Valid/", GzipCodec.class);
inputDataRDD
.filter(filterData -> filterData.split("\t").length != columnLength)
.coalesce(1)
.saveAsTextFile(args[2]+"Error/", GzipCodec.class);
片段结束..
该文件中有近 10M 条记录。
Java 和 Python 之间的 sc.textfile (filename , int numPartitions) 有什么不同吗?还是我遗漏了什么。
需要你的帮助来找出我犯的错误。
注意:-使用 eclipse 构建了一个 maven 并在 Yarn 中运行了以下命令。
spark-submit --class com.virtualpairprogrammers.ProcessFilesToHDFS --master yarn learningSpark-0.0.1-SNAPSHOT.jar "/input/ABFeeds/" "ABFeeds_2020-04-20.tsv.gz" "/output/ABFeeds/2020-05-06/" 996
提前致谢
问候
山姆
【问题讨论】:
-
当有效记录被发送到错误文件夹时,您是否检查了 Spark Java 程序正在计算的列数是多少?
-
@kaysush - 传递的列数是有效和错误的常数。两者之间的任何时候都不会改变。
-
抱歉,我表达的不够清楚。我想知道
filterData.split("\t").length表达式返回的有效行是什么? -
@kaysush:我尝试使用以下语句打印以下内容,但运行时未打印。 inputDataRDD.foreach(value -> System.out.println(value.split("\t").length));
-
@kaysush:我想我知道为什么它不打印了。为我的示例数据修复了它。让我为我的有效数据运行它,看看输出是什么。
标签: java apache-spark rdd