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
相关产品推荐
相关产品推荐

