PySpark读取MariaDB数据报错:日期解码失败与Schema参数错误
本地Windows环境下安装了MariaDB和PySpark 3.5.1,尝试读取MariaDB中stock_prices_nyse表数据到DataFrame时遇到两类错误:
错误1:不指定Schema读取时的日期解码失败
Py4JJavaError: An error occurred while calling o86.showString.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 3.0 failed 1 times, most recent failure: Lost task 0.0 in stage 3.0 (TID 3) (WinPc24.ht.home executor driver): java.sql.SQLDataException: value 'date' (VARSTRING) cannot be decoded as Date
at org.mariadb.jdbc.client.column.StringColumn.decodeDateText(StringColumn.java:307)
错误2:指定Schema读取时参数不支持
TypeError: DataFrameReader.jdbc() got an unexpected keyword argument 'schema'
MariaDB表结构
DESC stock_prices_nyse; +---------+----------------+------+-----+---------+-------+ | Field | Type | Null | Key | Default | Extra | +---------+----------------+------+-----+---------+-------+ | date | date | YES | | NULL | | | symbol | varchar(10) | YES | | NULL | | | open | decimal(15,10) | YES | | NULL | | | close | decimal(15,10) | YES | | NULL | | | low | decimal(15,10) | YES | | NULL | | | high | decimal(15,10) | YES | | NULL | | | volume | int(11) | YES | | NULL | | | country | char(10) | YES | | NULL | | +---------+----------------+------+-----+---------+-------+
原始测试代码
from pyspark.sql import SparkSession jdbc_driver_path = "C:/Spark/spark-3.5.1-bin-hadoop3/jars/mariadb-java-client-3.4.0.jar" spark = ( SparkSession.builder .appName("App2-NYSE") .config("spark.driver.extraClassPath", jdbc_driver_path) .getOrCreate() ) # Define the database properties db_properties = { "driver": "org.mariadb.jdbc.Driver", "url": "jdbc:mariadb://localhost:3306/name-database1", "user": "root", "password": "pass-mariadb" } table_name = "stock_prices_nyse" # Define the schema of the dataframe based on the table of MariaDB from pyspark.sql.types import StructType, StructField, DateType, StringType, DecimalType, IntegerType schema = StructType([ StructField("date", DateType(), True), StructField("symbol", StringType(), True), StructField("open", DecimalType(15, 10), True), StructField("close", DecimalType(15, 10), True), StructField("low", DecimalType(15, 10), True), StructField("high", DecimalType(15, 10), True), StructField("volume", IntegerType(), True), StructField("country", StringType(), True) ]) # Read data from the table with the specified schema df = spark.read.jdbc(url=db_properties["url"], table=table_name, properties=db_properties, schema=schema) print("First 5 rows of the dataframe:\n", df.show(5)) # Read data from the table without schema df = spark.read.jdbc(url=db_properties["url"], table=table_name, properties=db_properties) print("First 5 rows of the dataframe:\n", df.show(5))
解决方案
针对错误1(日期解码失败)
错误根源是表中date字段存在无法解析为日期的字符串值(比如字符串date),或JDBC驱动与日期格式不兼容。解决方式:
- 清理脏数据:检查MariaDB表,找到
date字段值为非日期格式的记录(如值为date的行),删除或修正这些数据。 - 调整JDBC URL参数:在数据库URL中添加时区和日期处理参数,强制驱动正确解析日期类型:
db_properties["url"] = "jdbc:mariadb://localhost:3306/name-database1?useLegacyDatetimeCode=false&serverTimezone=UTC" - 先按字符串读取再转换:如果无法清理数据,先读取所有字段为字符串,再将
date列转换为Date类型:# 先读取为原始类型 df_raw = spark.read.jdbc(url=db_properties["url"], table=table_name, properties=db_properties) # 转换date列类型并过滤无效值 from pyspark.sql.functions import to_date, col df = df_raw.withColumn("date", to_date(col("date"), "yyyy-MM-dd")) \ .filter(col("date").isNotNull())
针对错误2(schema参数不支持)
PySpark的jdbc()方法不支持直接传入schema参数,schema仅适用于文件类数据源。替代方案:
- 读取后转换字段类型:先读取数据,再通过
cast方法将字段转换为目标类型:df = spark.read.jdbc(url=db_properties["url"], table=table_name, properties=db_properties) df = df.withColumn("date", col("date").cast(DateType())) \ .withColumn("open", col("open").cast(DecimalType(15,10))) \ .withColumn("close", col("close").cast(DecimalType(15,10))) \ .withColumn("low", col("low").cast(DecimalType(15,10))) \ .withColumn("high", col("high").cast(DecimalType(15,10))) \ .withColumn("volume", col("volume").cast(IntegerType())) - 使用自定义SQL查询:通过子查询包装表读取数据,避免直接读取表时的类型识别问题:
# 用子查询包装表,避免JDBC直接读取的类型映射问题 query = "(SELECT date, symbol, open, close, low, high, volume, country FROM stock_prices_nyse) AS tmp_table" df = spark.read.jdbc(url=db_properties["url"], table=query, properties=db_properties) # 按需转换字段类型
内容的提问来源于stack exchange,提问作者ALIMUL HASAN

