【问题标题】:split the file into multiple files based on a string in spark scala根据 spark scala 中的字符串将文件拆分为多个文件
【发布时间】:2018-02-24 12:40:29
【问题描述】:

我有一个文本文件,其以下数据没有特定格式

abc*123     *180109*1005*^*001*0000001*0*T*:~
efg*05*1*X*005010X2A1~
k7*IT 1234*P*234df~ 
hig*0109*10052200*Rq~
abc*234*9698*709870*99999*N:~
tng****MI*917937861~
k7*IT 8876*e*278df~
dtp*D8*20171015~

我希望输出如下两个文件:

基于字符串abc,我想拆分文件。

文件 1:

abc*123     *180109*1005*^*001*0000001*0*T*:~
efg*05*1*X*005010X2A1~
k7*IT 1234*P*234df~ 
hig*0109*10052200*Rq~

文件 2:

abc*234*9698*709870*99999*N:~
tng****MI*917937861~
k7*IT 8876*e*278df~
dtp*D8*20171015~

并且文件名应该是 IT 名称(行以 k7 开头)所以 file1 名称应该是 IT_1234 第二个文件名应该是 IT_8876。

【问题讨论】:

  • 每个新文件会有多少数据?
  • 为什么要用 Spark 做这个?为什么不说,bash?文件在 HDFS 上吗?
  • 是的,您能告诉我们这样做的最终目标吗?将帮助我们找到合适的解决方案。
  • 是的,我们每小时获得 2.5 GB 的数据,文件位于 HDFS 中。

标签: scala apache-spark


【解决方案1】:

我在一个项目中使用了这个肮脏的小技巧:

sc.hadoopConfiguration.set("textinputformat.record.delimiter", "abc")

您可以设置用于读取文件的 spark 上下文的分隔符。所以你可以做这样的事情:

val delimit = "abc"
sc.hadoopConfiguration.set("textinputformat.record.delimiter", delimit)
val df = sc.textFile("your_original_file.txt")
           .map(x => (delimit ++ x))
           .toDF("delimit_column")
           .filter(col("delimit_column") !== delimit)

然后您可以将要写入的 DataFrame(或 RDD)的每个元素映射到文件中。

这是一个肮脏的方法,但它可能会帮助你!

祝你有美好的一天

PS:最后的过滤器是删除带有连接分隔符的第一行

【讨论】:

  • 谢谢。我可以使用这个解决方案。我已经编辑了这个问题,请帮我解决这个问题
  • 您好,我认为您不能在 Spark 中命名文件。您应该在编写之前使用 hadoop 库并使用文件名创建路径。或者创建一个shell脚本
  • 嗨,如果我在记录中间有“abc”字符串,在这种情况下,而不是 2,我将在输出中得到 3 条记录。我们可以在定义时使用子字符串函数吗?定界
【解决方案2】:

您可以受益于 sparkContext 的 wholeTextFiles 函数来读取文件。然后解析它来分隔字符串(这里我使用####作为不会在文本中重复的字符的独特组合)

val rdd = sc.wholeTextFiles("path to the file")
  .flatMap(tuple => tuple._2.replace("\r\nabc", "####abc").split("####")).collect()

然后循环数组以将文本保存到输出

for(str <- rdd){
  //saving codes here
}

【讨论】:

    猜你喜欢
    • 2017-08-28
    • 1970-01-01
    • 2012-07-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-13
    相关资源
    最近更新 更多