【发布时间】:2021-03-08 22:28:50
【问题描述】:
我正在尝试从 sqlserver 读取表并在读取时应用分区。在读取数据之前,我想得到下界和上界的界限,如下所示。
boundsDF = spark.read.format('jdbc')
.option('url', 'url')
.option('driver', 'com.microsoft.sqlserver.jdbc.SQLServerDriver')
.option('user', username)
.option('password', password)
.option('dbtable', f'(select min(updated_datetime) as mint, max(updated_datetime) as maxt from tablename)
.load()
我从 boundsDF 中提取了如下值:
maxdate = [x["maxt"] for x in boundsDF.rdd.collect()]
mindate = [x["mint"] for x in boundsDF.rdd.collect()]
这就是我在阅读时指定时间戳列的方式:
dataframe = spark.read.format('jdbc')
.option('url', url)
.option('driver', 'com.microsoft.sqlserver.jdbc.SQLServerDriver')
.option('user', user)
.option('password', password)
.option('dbtable', tablename)
.option('partitionColumn', timestamp_column)
.option('numPartitions', 3)
.option('lowerBound', mindate[0])
.option('upperBound', maxdate[0])
.option('fetchsize', 5000)
.load()
如果我打印下面的 mindate 和 maxdate 的值,它们的外观如下:
mindate[0]: datetime.datetime(2010, 10, 4, 11, 54, 13, 543000)
maxdate[0]: datetime.datetime(2021, 3, 5, 17, 59, 45, 880000)
当我打印dataframe.count() 时,我看到如下异常消息。
例外:
org.apache.spark.SparkException: Job aborted due to stage failure: Task 2 in stage 18.0 failed 1 times, most recent failure: Lost task 2.0 in stage 18.0 (TID 21, executor driver): com.microsoft.sqlserver.jdbc.SQLServerException: Conversion failed when converting date and/or time from character string.
自从我开始使用 Spark 以来,我一直使用整数列作为我的分区列。这是我第一次使用时间戳列对数据进行分区。
mindate[0] 和 maxdate[0] 的格式是否正确,可以在我的读取语句中指定? 如果我以正确的方式实现代码,谁能告诉我?
【问题讨论】:
-
您必须使用 SQL Server 可以理解的格式将参数作为字符串传递
-
但我看到一条错误消息,说无法理解边界传递的格式。所以我想我可以这样尝试。
-
但该语法在 Oracle 中。我正在使用 SqlServer。
标签: sql-server apache-spark pyspark apache-spark-sql