【问题标题】:load big japanese files in hadoop在hadoop中加载大日本文件
【发布时间】:2017-09-18 01:42:16
【问题描述】:

我在 HDFS 上有这个巨大的文件,它是我的数据库的提取。例如:

1||||||1||||||||||||||0002||01||1999-06-01 16:18:38||||2999-12-31 00:00:00||||||||||||||||||||||||||||||||||||||||||||||||||||||||2||||0||W.ISHIHARA||||1999-06-01 16:18:38||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||19155||||||||||||||1||1||NBV||||||||||||||U||||||||N||||||||||||||||||||||
1||||||8||2000-08-25 00:00:00||||||||3||||0001||01||1999-06-01 16:26:16||||1999-06-01 17:57:10||||||||||300||||||PH||400||Yes||PH||0255097�`||400||||1||103520||||||1||4||10||||20||||||||||2||||0||S.OSARI||1961-10-05 00:00:00||1999-06-01 16:26:16||�o��������������||�o��������������||1||||����||||1||1994-01-24 00:00:00||2||||||75||1999-08-25 00:00:00||1999-08-25 00:00:00||0||1||||4||||||�l��������������||�o��������������||�l��������������||||�o��������������||NP||||�l��������������||�l��������������||||||5||19055||||||||||1||||8||1||NBV||||||||||||||U||||||||N||||||||||||||||||||||
  • 文件大小:40GB
  • 记录数:~120 000 000
  • 字段数:112
  • 字段分隔:||
  • 行 sep:\n
  • 编码:sjis

我想使用 pyspark(1.6 和 python 3)在 hive 中加载这个文件。但我的工作一直失败。 这是我的代码:

toProcessFileDF = sc.binaryFiles("MyFile")\
    .flatMap(lambda x: x[1].split(b'\n'))\
    .map(lambda x: x.decode('sjis'))\
    .filter(lambda x: x.count('|')==sepCnt*2)\
    .map(lambda x: x.split('||'))\
    .toDF(schema=tableSchema) #tableSchema is the schema retrieved from hive
toProcessFileDF.write.saveAsTable(tableName, mode='append')

我收到了几个错误,其中包括 jave 143(内存错误)、心跳超时和内核已死。 (如果您需要确切的日志错误,请告诉我)。

这是正确的方法吗?也许有更聪明或更有效的方法。您能告诉我如何执行此操作吗?

【问题讨论】:

  • 您是否尝试在文件上简单地创建一个外部配置单元表?我认为 Spark 不会对小型集群上 40 GB 的内存数据感到高兴
  • 2 个问题。我的架构中的 Hive 对文件所在的文件夹没有读取权限。 Hive 外部表不支持 2 字符字段 sep。但我的集群并不小。
  • 移动或复制 HDFS 文件,然后呢?我可能错了,但FIELDS TERMINATED BY '||' 似乎对我有用

标签: python apache-spark hive pyspark hdfs


【解决方案1】:

我发现 databrick csv 阅读器非常有用。

toProcessFileDF_raw = sqlContext.read.format('com.databricks.spark.csv')\
                                        .options(header='false',
                                                 inferschema='false',
                                                 charset='shift-jis',
                                                 delimiter='\t')\
                                        .load(toProcessFile)

不幸的是,我可以使用分隔符选项仅使用一个字符进行拆分。因此,我的解决方案是使用制表符拆分,因为我确定我的文件中没有任何内容。然后我可以在我的线上应用拆分。

这并不完美,但至少我有正确的编码并且我不会将所有内容都放在内存中。

【讨论】:

    【解决方案2】:

    日志会很有帮助。

    binaryFiles 会将这个二进制文件从 HDFS 读取为单个记录,并以键值对的形式返回,其中键是每个文件的路径,值是每个文件的内容。

    Spark 文档中的注释:首选小文件,也允许使用大文件,但可能会导致性能不佳。

    如果你使用textFile会更好

    toProcessFileDF = sc.textFile("MyFile")\
                        .map(lambda line: line.split("||"))\
    ....
    

    另外,一种选择是在读取初始文本文件时指定更多的分区(例如 sc.textFile(path, 200000)),而不是在读取后重新分区。

    另一件重要的事情是确保您的输入文件是可拆分的(某些压缩选项使其不可拆分,在这种情况下,Spark 可能必须在单台机器上读取它,从而导致 OOM)。

    【讨论】:

    • 这绝对是朝方向迈出的一步,但考虑到问题你会像sc.textFile(path, use_unicode=False).map(lambda s: s.decode("sjis"))
    • sc.textFile 接收 UTF-8 格式的文本......因此,之后无法解码。 utf-8 中不可读的字符已损坏。
    • @zero323 仅供参考,这个想法很好,但“use_unicode”并没有像你想象的那样工作。它带来了 utf8 解码文件而不是 utf8 编码文件......这会在日语字符中产生错误。
    猜你喜欢
    • 2015-11-27
    • 1970-01-01
    • 1970-01-01
    • 2016-01-21
    • 2016-01-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-28
    相关资源
    最近更新 更多