Apache Spark读取JSON调用TMDB API时列解析错误求助
PySpark脚本填充
produtora列时出现UNRESOLVED_COLUMN错误的解决方法 问题概述
编写PySpark脚本读取文件夹中的JSON文件,调用TMDB API获取制作公司数据填充到produtora列,运行时抛出UNRESOLVED_COLUMN错误,提示produtora列无法解析。
错误栈信息
Traceback (most recent call last): File "/home/gwillye/Documentos/Lab Ex/dadosProdutora.py", line 41, in <module> dados_json = dados_json.withColumn("produtora", col("produtora").cast(StringType())) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/gwillye/spark/python/lib/pyspark.zip/pyspark/sql/dataframe.py", line 5170, in withColumn File "/home/gwillye/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in __call__ File "/home/gwillye/spark/python/lib/pyspark.zip/pyspark/errors/exceptions/captured.py", line 185, in deco pyspark.errors.exceptions.captured.AnalysisException: [UNRESOLVED_COLUMN.WITH_SUGGESTION] A column or function parameter with name `produtora` cannot be resolved. Did you mean one of the following? [`genero`, `id`, `profissao`, `notaMedia`, `numeroVotos`].; 'Project [_corrupt_record#8, anoFalecimento#9, anoLancamento#10L, anoNascimento#11, anoTermino#12, genero#13, generoArtista#14, id#15, nomeArtista#16, notaMedia#17, numeroVotos#18L, personagem#19, profissao#20, tempoMinutos#21, tituloOriginal#22, tituloPincipal#23, titulosMaisConhecidos#24, cast('produtora as string) AS produtora#42] +- Relation [_corrupt_record#8,anoFalecimento#9,anoLancamento#10L,anoNascimento#11,anoTermino#12,genero#13,generoArtista#14,id#15,nomeArtista#16,notaMedia#17,numeroVotos#18L,personagem#19,profissao#20,tempoMinutos#21,tituloOriginal#22,tituloPincipal#23,titulosMaisConhecidos#24] json'''
JSON文件示例
[ { "id": "tt0061797", "tituloPincipal": "The Madcap Island", "tituloOriginal": "Hyokkori hyōtan-jima", "anoLancamento": 1967.0, "tempoMinutos": 61.0, "genero": "Animation", "notaMedia": 5.6, "numeroVotos": 10, "generoArtista": "actress", "personagem": null, "nomeArtista": "Chinatsu Nakayama", "anoNascimento": 1948.0, "anoFalecimento": null, "profissao": "actress,soundtrack,music_department", "titulosMaisConhecidos": "tt1734449,tt0204339,tt0202407,tt0081881", "produtora": null }, { "id": "tt0061797", "tituloPincipal": "The Madcap Island", "tituloOriginal": "Hyokkori hyōtan-jima", "anoLancamento": 1967.0, "tempoMinutos": 61.0, "genero": "Animation", "notaMedia": 5.6, "numeroVotos": 10, "generoArtista": "actor", "personagem": null, "nomeArtista": "Arihiro Fujimura", "anoNascimento": 1934.0, "anoFalecimento": 1982.0, "profissao": "actor,soundtrack", "titulosMaisConhecidos": "tt0060926,tt0102587,tt0298290,tt0062224", "produtora": null } ]
原脚本代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf from pyspark.sql.types import StringType import requests def get_tmdb_data(movie_id): api_key = 'f905807b2900febaccfb008c16388168' url = f'https://api.themoviedb.org/3/movie/{movie_id}?api_key={api_key}&language=pt-BR' response = requests.get(url) if response.status_code == 200: return response.json() else: return None def get_produtora_udf(movie_id): tmdb_data = get_tmdb_data(movie_id) if tmdb_data and 'production_companies' in tmdb_data: produtoras = tmdb_data['production_companies'] if produtoras: for produtora in produtoras: if 'name' in produtora: return produtora['name'] return None def processar_registro(row): movie_id = row['id'] produtora_atual = row['produtora'] if produtora_atual is None: nova_produtora = get_produtora_udf(movie_id) return row.asDict() | {'produtora': nova_produtora} else: return row if __name__ == "__main__": spark = SparkSession.builder.appName("ExemploPyspark").getOrCreate() # Leitura dos arquivos JSON na pasta 'JSON/' dados_json = spark.read.json("JSON/") # Adiciona a coluna 'produtora' com valor null dados_json = dados_json.withColumn("produtora", col("produtora").cast(StringType())) # Aplica a função processar_registro para substituir 'null' nas colunas 'produtora' dados_processados = dados_json.rdd.map(processar_registro).toDF() # Exibe o DataFrame resultante dados_processados.show(truncate=False) # Encerra a sessão Spark spark.stop()
问题根源
Spark自动推断JSON schema时,若所有记录的produtora字段值为null,或者部分记录缺失该字段,Spark会忽略该列,导致读取后的DataFrame中不存在produtora列。脚本中直接尝试引用col("produtora")进行类型转换,触发列不存在的错误。
解决方案
1. 确保produtora列存在
读取JSON后先检查列是否存在,若不存在则添加值为null的produtora列,再统一转换为StringType。
2. 使用Spark UDF替代RDD转换
避免直接操作RDD,改用DataFrame API的UDF,利用Spark的优化机制提升性能。
修改后的完整脚本
from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf, lit, when from pyspark.sql.types import StringType import requests def get_tmdb_data(movie_id): api_key = 'f905807b2900febaccfb008c16388168' # TMDB电影ID为纯数字,需去除原ID的'tt'前缀 tmdb_movie_id = movie_id.lstrip('tt') url = f'https://api.themoviedb.org/3/movie/{tmdb_movie_id}?api_key={api_key}&language=pt-BR' response = requests.get(url) if response.status_code == 200: return response.json() else: return None def get_produtora(movie_id): tmdb_data = get_tmdb_data(movie_id) if tmdb_data and 'production_companies' in tmdb_data: produtoras = tmdb_data['production_companies'] if produtoras: for produtora in produtoras: if 'name' in produtora: return produtora['name'] return None # 注册UDF get_produtora_udf = udf(get_produtora, StringType()) if __name__ == "__main__": spark = SparkSession.builder.appName("ExemploPyspark").getOrCreate() # 读取JSON文件 dados_json = spark.read.json("JSON/") # 确保produtora列存在并转换为StringType if "produtora" not in dados_json.columns: dados_json = dados_json.withColumn("produtora", lit(None).cast(StringType())) else: dados_json = dados_json.withColumn("produtora", col("produtora").cast(StringType())) # 使用UDF填充null值的produtora列 dados_processados = dados_json.withColumn( "produtora", when(col("produtora").isNull(), get_produtora_udf(col("id"))).otherwise(col("produtora")) ) # 显示结果 dados_processados.show(truncate=False) # 停止Spark会话 spark.stop()
额外说明
- TMDB API的电影ID是纯数字,原JSON中的
id带tt前缀,需去除后再调用API,否则会返回404错误; - 频繁调用TMDB API可能触发限流,建议添加请求延迟、缓存已获取的数据,或使用批量查询接口优化性能。
内容的提问来源于stack exchange,提问作者Gabriel Willye Borges
相关产品推荐
相关产品推荐

