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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 16:32:54