You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kedro读取Spark.SparkDataSet:如何定义列名及解决Schema报错

在Kedro中读取SparkDataSet时如何定义列名?

问题说明

我当前的catalog.yaml配置如下:

user-playlists: 
  type: spark.SparkDataSet
  file_format: csv
  filepath: data/01_raw/lastfm-dataset-1K/userid-timestamp-artid-artname-traid-traname.tsv
  load_args:
    sep: "\t"
    header: False
#    schema:
#      filepath: conf/base/playlists-schema.json
  save_args:
    index: False

尝试用以下JSON Schema定义列名时,系统报错:schema Please provide a valid JSON-serialised 'pyspark.sql.types.StructType'.

{
  "fields": [
    {"name": "userid", "type": "string", "nullable": true},
    {"name": "timestamp", "type": "string", "nullable": true},
    {"name": "artid", "type": "string", "nullable": true},
    {"name": "artname", "type": "string", "nullable": true},
    {"name": "traid", "type": "string", "nullable": true},
    {"name": "traname", "type": "string", "nullable": true}
  ],
  "type": "struct"
}

解决办法

方法1:修正JSON Schema格式

PySpark对序列化的StructType格式有严格要求,必须符合Spark的类型定义规范。你可以把Schema改成以下两种格式之一:

格式一(带类型嵌套)

{
  "type": "struct",
  "fields": [
    {"name": "userid", "type": {"type": "string"}, "nullable": true},
    {"name": "timestamp", "type": {"type": "string"}, "nullable": true},
    {"name": "artid", "type": {"type": "string"}, "nullable": true},
    {"name": "artname", "type": {"type": "string"}, "nullable": true},
    {"name": "traid", "type": {"type": "string"}, "nullable": true},
    {"name": "traname", "type": {"type": "string"}, "nullable": true}
  ]
}

格式二(带metadata字段)

{
  "fields": [
    {"metadata": {}, "name": "userid", "nullable": true, "type": "string"},
    {"metadata": {}, "name": "timestamp", "nullable": true, "type": "string"},
    {"metadata": {}, "name": "artid", "nullable": true, "type": "string"},
    {"metadata": {}, "name": "artname", "nullable": true, "type": "string"},
    {"metadata": {}, "name": "traid", "nullable": true, "type": "string"},
    {"metadata": {}, "name": "traname", "nullable": true, "type": "string"}
  ],
  "type": "struct"
}

修改完成后,取消catalog.yaml中schema配置的注释,确保文件路径指向正确的Schema文件即可。

方法2:直接在load_args中指定列名

如果不需要严格约束数据类型,这是最简洁的方式——直接在load_args里添加names参数,把列名按顺序列出来:

user-playlists: 
  type: spark.SparkDataSet
  file_format: csv
  filepath: data/01_raw/lastfm-dataset-1K/userid-timestamp-artid-artname-traid-traname.tsv
  load_args:
    sep: "\t"
    header: False
    names: ["userid", "timestamp", "artid", "artname", "traid", "traname"]
  save_args:
    index: False

方法3:用代码定义Schema(灵活控制类型)

如果需要更精细的类型控制,比如后续要修改类型、添加约束,可以在Kedro项目的hooks.py或自定义模块里用PySpark代码定义Schema:

from pyspark.sql.types import StructType, StructField, StringType

playlist_schema = StructType([
    StructField("userid", StringType(), nullable=True),
    StructField("timestamp", StringType(), nullable=True),
    StructField("artid", StringType(), nullable=True),
    StructField("artname", StringType(), nullable=True),
    StructField("traid", StringType(), nullable=True),
    StructField("traname", StringType(), nullable=True)
])

之后在catalog.yaml中通过引用该变量的方式关联Schema(需要确保代码能被Kedro的上下文加载到)。


内容的提问来源于stack exchange,提问作者gaut

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 02:07:05