【问题标题】:How to specify timestamp values for bounds in pyspark dataframe?如何为 pyspark 数据框中的边界指定时间戳值?
【发布时间】: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


【解决方案1】:

问题是您在 SQL 表中使用什么数据类型?

  1. TIMESTAMP 不是日期时间数据类型。它是一个内部行版本号(二进制),与时间数据无关
  2. DATETIME 是旧的不推荐使用的 DATE + TIME 数据类型,秒数限制为 3 位小数
  3. DATETIME2 取代了日期时间,是用于 DATE + TIME 的新数据类型,并且有一个限制,您可以在 0 到 7 位小数之间选择秒

现在有两句话:

  1. 如果您使用 TIMESTAMP,请将其替换为 DATETIME2 所需的精度(默认为 7)。
  2. 如果您使用 DATETIME 并且不想将其替换为 DATETIME2,则只能为秒的小数部分指定 3 位数字,但我在您的代码中看到的是 mindate[0]: datetime.datetime(2010 , 10, 4, 11, 54, 13, 543000)

DATETIME2 比限制为 3 毫秒的 DATETIME 更准确,导致某些查询被错误解释

【讨论】:

  • 我理解你的解释。但是当我给数据时间值是我的界限时,我面临着我在问题中提到的异常。
  • 就像我说的,如果你使用 DATETIME 并且你不想用 DATETIME2 替换它,你只能指定 3 位作为秒的小数部分,所以修改这个代码:mindate[0] : datetime.datetime(2010, 10, 4, 11, 54, 13, 543000),如果你使用 DATETIME 并且你不想用 datetime.datetime(2010, 10, 4, 11, 54, 13,第543章)
  • 关于备注部分中您的答案第 2 点,Spark 默认以 datetime.datetime 格式读取该日期列。您建议我将 Datetime 的现有格式更改为 Datetime2。如果我必须将其转换为 Datetime2,我还应该在通过 Spark 读取数据时将整个列转换为 Datetime2,因为我以该特定格式给出了分区边界。我的理解正确吗?
  • 我想是的,但你必须尝试
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-09-11
  • 2022-01-18
  • 1970-01-01
  • 2021-09-06
  • 2021-02-06
  • 2012-12-09
相关资源
最近更新 更多