DataFrame - ValueError:StructType 的意外元组

fly*_*zai 4 python dataframe apache-spark apache-spark-sql pyspark

我正在尝试为dataframe. 我传入的数据是从json. 这是我的初始数据:

json2 = sc.parallelize(['{"name": "mission", "pandas": {"attributes": "[0.4, 0.5]", "pt": "giant", "id": "1", "zip": "94110", "happy": "True"}}'])
Run Code Online (Sandbox Code Playgroud)

然后这里是如何指定架构:

schema = StructType(fields=[
    StructField(
        name='name',
        dataType=StringType(),
        nullable=True
    ),
    StructField(
        name='pandas',
        dataType=ArrayType(
            StructType(
                fields=[
                    StructField(
                        name='id',
                        dataType=StringType(),
                        nullable=False
                    ),
                    StructField(
                        name='zip',
                        dataType=StringType(),
                        nullable=True
                    ),
                    StructField(
                        name='pt',
                        dataType=StringType(),
                        nullable=True
                    ),
                    StructField(
                        name='happy',
                        dataType=BooleanType(),
                        nullable=False
                    ),
                    StructField(
                        name='attributes',
                        dataType=ArrayType(
                            elementType=DoubleType(),
                            containsNull=False
                        ),
                        nullable=True

                    )
                ]
            ),
            containsNull=True
        ),
        nullable=True
    )
])
Run Code Online (Sandbox Code Playgroud)

当我使用sqlContext.createDataFrame(json2, schema)然后尝试对结果执行操作show()时,dataframe我收到以下错误:

ValueError: Unexpected tuple '{"name": "mission", "pandas": {"attributes": "[0.4, 0.5]", "pt": "giant", "id": "1", "zip": "94110", "happy": "True"}}' with StructType
Run Code Online (Sandbox Code Playgroud)

zer*_*323 5

首先json2只是一个RDD[String]. Spark 没有关于用于编码数据的序列化格式的特殊知识。此外,它期望一个RDD或Row多个产品,但显然并非如此。

在 Scala 中,您可以使用

sqlContext.read.schema(schema).json(rdd) 
Run Code Online (Sandbox Code Playgroud)

有RDD[String],但有两个问题:

  • 这种方法在 PySpark 中不能直接访问
  • 即使它是您创建的架构也是无效的:

    • pandas是struct不和array
    • pandas.happy不是string一个boolean
    • pandas.attributes是string不是array

Schema 仅用于避免类型推断,而不用于类型转换或任何其他转换。如果你想转换数据,你必须先解析它:

def parse(s: str) -> Row:
    return ...

rdd.map(parse).toDF(schema)
Run Code Online (Sandbox Code Playgroud)

假设你有这样的 JSON(固定类型):

{"name": "mission", "pandas": {"attributes": [0.4, 0.5], "pt": "giant", "id": "1", "zip": "94110", "happy": true}} 
Run Code Online (Sandbox Code Playgroud)

正确的架构如下所示

StructType([
    StructField("name", StringType(), True),
    StructField("pandas", StructType([
        StructField("attributes", ArrayType(DoubleType(), True), True),
        StructField("happy", BooleanType(), True),
        StructField("id", StringType(), True),
        StructField("pt", StringType(), True),
        StructField("zip", StringType(), True))],
    True)])
Run Code Online (Sandbox Code Playgroud)