PySpark执行PostgreSQL INSERT/UPDATE语句报错求助
问题原因及解决方法
报错根源
你使用的spark.read.jdbc()是Spark专为读取数据、生成DataFrame设计的API,仅支持执行能返回结果集的SQL语句(如SELECT)。而INSERT/UPDATE这类DML语句不会返回结果集,用read方式执行时,PostgreSQL驱动会将其当作查询语句解析,直接触发语法错误。
正确实现方案
既然不能用psycopg2/jaydebeapi,可通过PySpark调用底层Java JDBC API执行DML语句,具体代码如下:
from pyspark.sql import SparkSession def run_pg_sql_spark(qry): url = "jdbc:postgresql://account-name" user = "user" password = "pwd" driver = "org.postgresql.Driver" # 初始化SparkSession(避免重复创建) spark = SparkSession.builder.appName("pg_db_conn").getOrCreate() # 获取JVM中的JDBC工具类 jvm = spark.sparkContext._jvm jvm.Class.forName(driver) # 建立连接并执行SQL conn = jvm.java.sql.DriverManager.getConnection(url, user, password) stmt = conn.createStatement() try: # executeUpdate专门处理无结果集的DML语句,返回受影响行数 affected_rows = stmt.executeUpdate(qry) print(f"SQL执行成功,受影响行数:{affected_rows}") finally: # 强制关闭资源,避免连接泄漏 stmt.close() conn.close() # 调用测试 run_pg_sql_spark("INSERT INTO table (task_name) values('x')")
注意事项
- 确保PostgreSQL JDBC驱动已添加到Spark的classpath中(比如启动Spark时用
--jars参数指定驱动包路径) executeUpdate()方法适用于所有不返回结果集的SQL操作(INSERT/UPDATE/DELETE/DDL等)
内容的提问来源于stack exchange,提问作者AnumNuma
相关产品推荐
相关产品推荐

