【问题标题】:Read a bytes column in spark在 spark 中读取字节列
【发布时间】:2018-08-21 22:53:14
【问题描述】:

我有一个数据集,其中包含一个未知(且不友好)编码的 ID 字段。我可以使用普通 python 读取单个列,并验证多个数据集的值是否不同且一致(即它可以用作连接的主键)。

使用spark.read.csv 加载文件时,spark 似乎正在将该列转换为utf-8。但是,一些多字节序列被转换为 Unicode 字符 U+FFFD REPLACEMENT CHARACTER.EF BF BD 十六进制)。

有没有办法强制 Spark 将列读取为字节而不是字符串?

这是一些可用于重新创建我的问题的代码(让列 a 成为 ID 字段):

使用示例数据创建文件

data = [
    (bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb0'), '1', 'a'),
    (bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb1'), '2', 'b'),
    (bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb2'), '3', 'c')
]

with open('sample.csv', 'wb') as f:
    header = ["a", "b", "c"]
    f.write(",".join(header)+"\n")
    for d in data:
        f.write(",".join(d) + "\n")

使用 Pandas 阅读

import pandas as pd
df = pd.read_csv("sample.csv", converters={"a": lambda x: x.encode('hex')})
print(df)
#                  a  b  c
#0  baed858e91d4c7b0  1  a
#1  baed858e91d4c7b1  2  b
#2  baed858e91d4c7b2  3  c

尝试使用 Spark 读取同一个文件

spark_df = spark.read.csv("sample.csv", header=True)
spark_df.show()
#+-----+---+---+
#|a    |b  |c  |
#+-----+---+---+
#|�텎��ǰ|1  |a  |
#|�텎��DZ|2  |b  |
#|�텎��Dz|3  |c  |
#+-----+---+---+

哎呀!好的,那转换成hex怎么样?

import pyspark.sql.functions as f
spark_df.withColumn("a", f.hex("a")).show(truncate=False)
#+----------------------------+---+---+
#|a                           |b  |c  |
#+----------------------------+---+---+
#|EFBFBDED858EEFBFBDEFBFBDC7B0|1  |a  |
#|EFBFBDED858EEFBFBDEFBFBDC7B1|2  |b  |
#|EFBFBDED858EEFBFBDEFBFBDC7B2|3  |c  |
#+----------------------------+---+---+

(在此示例中,值是不同的,但在我的较大文件中并非如此)

如你所见,值是close,但部分字节已被EFBFBD替换

有没有办法在 Spark 中读取文件(可能使用rdd?),这样我的输出看起来就像熊猫版本:

#+----------------+---+---+
#|a               |b  |c  |
#+----------------+---+---+
#|baed858e91d4c7b0|1  |a  |
#|baed858e91d4c7b1|2  |b  |
#|baed858e91d4c7b2|3  |c  |
#+----------------+---+---+

我已尝试强制转换为 byte 并指定架构,以便此列为 ByteType(),但这不起作用。

编辑

我使用的是 Spark v 2.1。

【问题讨论】:

  • 您找到了解决上述问题的任何方法。我遇到了类似的问题
  • 能告诉我你找到的解决方案吗?

标签: apache-spark encoding pyspark apache-spark-sql


【解决方案1】:

问题的根源在于分隔文件不太适合二进制数据。

如果文本有已知的一致编码,请使用charset 选项。请参阅https://github.com/databricks/spark-csv#features(我不知道 2.x 文档中描述了分隔阅读选项的好地方,所以我仍然回到 1.x 文档)。我建议尝试使用 8 位 ASCII,例如 ISO-8859-1US-ASCII

如果没有这样的编码,您需要将输入转换为不同的格式,例如,对第一列进行 base64 编码,或者操作读取的数据以将其恢复为您需要的格式。

【讨论】:

  • 您能否提供一些适用于我帖子中示例的代码?
  • 不是真的,因为我没有数据并且需要进行实验。在使用 Spark 读取 CSV 时尝试传递 charset = "ISO-8859-1" 选项。
  • 我在我的问题中放入了生成数据的代码。感谢您的建议,但我认为这不会解决这个问题。
  • 你为什么不认为读取 8 位二进制(作为 ASCII)会起作用?这种方法基本上跳过了 Java 字符串处理代码的任何解码尝试。
  • @pault 您找到了解决上述问题的任何方法。我遇到了类似的问题
【解决方案2】:

如何将其存储为base 64编码并在读取时对其进行解码?

存储

import base64

data = [
    (base64.b64encode(bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb0')), '1', 'a'),
    (base64.b64encode(bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb1')), '2', 'b'),
    (base64.b64encode(bytes(b'\xba\xed\x85\x8e\x91\xd4\xc7\xb2')), '3', 'c')
]

with open('sample.csv', 'wb') as f:
    header = ["a", "b", "c"]
    f.write(",".join(header)+"\n")
    for d in data:
        f.write(",".join(d) + "\n")

阅读

import pyspark.sql.functions as f
import base64

spark_df.withColumn("a", base64.b64decode("a"))

【讨论】:

  • 重点是在spark中做到这一点。如果我可以修改源数据,我可以使用 pandas 重写它,如上所示。
【解决方案3】:

似乎正在进行一些 UTF-8 解码; \xba 不是任何有效的 UTF-8 编码(见下文),正在被 \uFFFD(“替换字符”)替换。这必然会发生:CSV 是一种文本格式,因此解码器必须假设某种编码来解释二进制文件。

让我们从第一个值 EFBFBD 开始并手动解码(https://en.wikipedia.org/wiki/UTF-8 可能有助于理解编码)。

  1. EF0b11101111,因此是一个从 1111 位开始的 3 字节序列(初始位 1110)。
  2. BF0b10111111,因此是位 111111 的延续(初始位 10)。
  3. BD0b10111101,所以继续使用位 111101。
  4. 将它们放在一起得到 nybbles 1111 1111 1111 1101,即 FFFD。

你不能用像 CSV 这样的文本格式来解决这个问题,除非你能让你的 CSV 编码器和解码器就某种二进制编码格式达成一致。该格式将是 binary,而不是 Unicode。请注意,如果你有一般的二进制值,你甚至不能相信行尾!

您不妨使用 base-64 编码。或者,如果您知道行尾不是问题(没有\x0a 和/或\x0c 字节),您可能可以使用行阅读器。但这可能不推荐。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-07-27
    • 2012-10-04
    • 1970-01-01
    • 1970-01-01
    • 2018-06-20
    • 2012-03-23
    • 1970-01-01
    • 2020-11-14
    相关资源
    最近更新 更多