PySpark执行CREATE TABLE遇mismatched input错误及字段失效问题
问题总结
- 执行
spark.sql('sql_query')报语法错误:mismatched input 'sql_query' expecting {EOF};去掉引号传入变量则触发NullPointerException - 修复上述错误后,通过
withColumn新增的type、timestamp字段未出现在目标Delta表中
错误原因与修复方案
问题1:Spark SQL执行错误
原因
- 用
spark.sql('sql_query')时,把字符串'sql_query'直接作为SQL语句传入,Spark会解析这个字符串本身,自然不符合SQL语法,抛出语法错误。 - 去掉引号后触发空指针,是因为
create_sql_statement函数没有返回生成的sql_query变量,调用函数后sql_query的值为None,传入spark.sql()导致空指针。
修复
修改create_sql_statement函数,最后添加返回语句:
def create_sql_statement(file_name): # 原函数代码不变... #return sql statement sql_query = f"CREATE TABLE {tbl_name} AS SELECT{sql_statement}FROM {best_tbl_match}" print(sql_query) return sql_query # 新增这一行
问题2:新增字段未生效
原因
spark.sql(sql_query)执行的是CREATE TABLE AS SELECT语句,这个语句的返回结果是DDL执行状态的DataFrame(通常为空或仅包含受影响行数),而不是查询出的表数据。后续调用withColumn是在这个空DF上操作,写入Delta表时自然不会包含新增字段。
修复
放弃用CREATE TABLE AS SELECT的方式,改为先查询数据、添加字段,再通过write.saveAsTable创建并写入表:
- 修改
create_sql_statement函数,返回查询数据的SQL语句(而非建表语句) - 在主逻辑中执行查询、添加字段后写入表
修改后的create_sql_statement函数:
def create_sql_statement(file_name): # 原函数代码不变... # 返回查询数据的SQL,而非建表语句 sql_query = f"SELECT{sql_statement}FROM {best_tbl_match}" print(sql_query) return sql_query
修改后的主逻辑代码:
file_list = [list of multiple files with the sql statements] for file in file_list: try: sql_query = create_sql_statement(file) # 执行查询获取数据DF,再添加字段 df = spark.sql(sql_query) \ .withColumn('type', F.lit('animal_type')) \ .withColumn('timestamp', F.current_timestamp()) # 写入时自动创建Delta表 df.write.format("delta") \ .option("overwriteSchema", "true") \ .mode("overwrite") \ .saveAsTable(f'{database}.{table.lower()}') except Exception as e: print(e)
额外优化建议
- 避免将Spark DataFrame转成Pandas处理SQL文本,直接用Spark API拼接字符串更高效:
# 替代原函数中读取文本并拼接的逻辑 one_string = df.select(F.concat_ws("", F.collect_list(df.value))).first()[0]
内容的提问来源于stack exchange,提问作者SunflowerParty
相关产品推荐
相关产品推荐

