如何在PySpark中自定义JDBC方言解决Cloud Spanner读写字符串字面量问题
问题描述
已成功通过JDBC将PySpark与Cloud Spanner建立连接,但在对表进行读写操作时遇到字符串字面量相关错误。经排查发现问题源于列名使用双引号("columns_name")与反引号(`columns_name`)的差异,现咨询如何在PySpark中自定义JDBC方言来解决该问题。
相关代码
from pyspark.sql import SparkSession from google.cloud import spanner from pyspark.sql.types import StructType, StructField, StringType, IntegerType import os OPERATION_TIMEOUT_SECONDS = 240 os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "./credentials_dev.json" credentials = "./credentials.json" spark = SparkSession.builder.appName( "Submit local CSV to Spanner").getOrCreate() spark.sparkContext.addFile("google-cloud-spanner-jdbc-2.8.0.jar") jdbc_url = "jdbc:cloudspanner:/projects/dev-data/" + \ "instances/spanner-test/databases/testdb?" + \ "credentials=credentials_dev.json;autocommit=false" driverClass = "com.google.cloud.spanner.jdbc.JdbcDriver" connection_properties = { 'driver': 'com.google.cloud.spanner.jdbc.JdbcDriver' } # ---------read operation----------------------------------------------- # df = spark.read.jdbc(url=jdbc_url,table='stories', properties=connection_properties) # print(df) # df.show() schema = StructType([ StructField("author", StringType(), True), StructField("s_by", StringType(), True), StructField("dead", StringType(), True), ]) csv = 'smalldata.csv' #reading local csv file-------- df1 = spark.read.schema(schema).csv("smalldata.csv",inferSchema=True,header=True) # df1.show() #write operation------------ df1.write.jdbc(url=jdbc_url, table="stories", mode="append",properties=connection_properties).save()
版本信息
- PySpark版本:3.3.2
- Cloud Spanner JDBC版本:2.8.0
解决方案
1. 自定义Cloud Spanner JDBC方言
PySpark默认JDBC方言会用反引号包裹列名,但Cloud Spanner JDBC驱动要求用双引号引用列名(尤其是列名含下划线等特殊字符时)。通过继承JdbcDialect类并实现指定方法,可修改这一行为:
from pyspark.sql.jdbc import JdbcDialect, registerDialect class CloudSpannerDialect(JdbcDialect): def canHandle(self, url: str) -> bool: # 匹配Cloud Spanner的JDBC URL前缀,确保方言仅对其生效 return url.startswith("jdbc:cloudspanner:") def quotedIdentifier(self, col_name: str) -> str: # 用双引号包裹列名,替代默认的反引号 return f'"{col_name}"' # 注册自定义方言 registerDialect(CloudSpannerDialect())
2. 修复代码冗余问题
原代码中df1.write.jdbc(...).save()存在冗余,jdbc方法本身已完成写入操作,无需额外调用save(),修改为:
# write operation------------ df1.write.jdbc(url=jdbc_url, table="stories", mode="append", properties=connection_properties)
3. 完整修改后的代码
from pyspark.sql import SparkSession from google.cloud import spanner from pyspark.sql.types import StructType, StructField, StringType, IntegerType from pyspark.sql.jdbc import JdbcDialect, registerDialect import os # 自定义Cloud Spanner JDBC方言 class CloudSpannerDialect(JdbcDialect): def canHandle(self, url: str) -> bool: return url.startswith("jdbc:cloudspanner:") def quotedIdentifier(self, col_name: str) -> str: return f'"{col_name}"' # 注册方言 registerDialect(CloudSpannerDialect()) OPERATION_TIMEOUT_SECONDS = 240 os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "./credentials_dev.json" credentials = "./credentials.json" spark = SparkSession.builder.appName( "Submit local CSV to Spanner").getOrCreate() spark.sparkContext.addFile("google-cloud-spanner-jdbc-2.8.0.jar") jdbc_url = "jdbc:cloudspanner:/projects/dev-data/" + \ "instances/spanner-test/databases/testdb?" + \ "credentials=credentials_dev.json;autocommit=false" driverClass = "com.google.cloud.spanner.jdbc.JdbcDriver" connection_properties = { 'driver': 'com.google.cloud.spanner.jdbc.JdbcDriver' } # ---------read operation----------------------------------------------- # df = spark.read.jdbc(url=jdbc_url,table='stories', properties=connection_properties) # print(df) # df.show() schema = StructType([ StructField("author", StringType(), True), StructField("s_by", StringType(), True), StructField("dead", StringType(), True), ]) csv = 'smalldata.csv' # reading local csv file-------- df1 = spark.read.schema(schema).csv("smalldata.csv", inferSchema=True, header=True) # df1.show() # write operation------------ df1.write.jdbc(url=jdbc_url, table="stories", mode="append", properties=connection_properties)
说明
canHandle方法负责识别Cloud Spanner的JDBC URL,避免自定义方言影响其他数据库连接。quotedIdentifier方法将列名用双引号包裹,完全适配Cloud Spanner的SQL语法要求。- 注册方言后,PySpark生成SQL语句时会自动使用双引号处理列名,彻底解决字面量格式不兼容问题。
内容的提问来源于stack exchange,提问作者Umesh Yadav
相关产品推荐
相关产品推荐

