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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 12:02:34