PySpark 毫秒的时间戳

Kee*_*pan 5 pyspark

我试图获得两个时间戳列之间的差异,但毫秒数消失了。

如何纠正这个?

from pyspark.sql.functions import unix_timestamp
timeFmt = "yyyy-MM-dd' 'HH:mm:ss.SSS"

data = [
    (1, '2018-07-25 17:15:06.39','2018-07-25 17:15:06.377'),
    (2,'2018-07-25 11:12:49.317','2018-07-25 11:12:48.883')

]

df = spark.createDataFrame(data, ['ID', 'max_ts','min_ts']).withColumn('diff',F.unix_timestamp('max_ts', format=timeFmt) - F.unix_timestamp('min_ts', format=timeFmt))
df.show(truncate = False)
Run Code Online (Sandbox Code Playgroud)

Tro*_*roy 8

假设您已经有一个包含时间戳类型列的数据框:

from datetime import datetime

data = [
    (1, datetime(2018, 7, 25, 17, 15, 6, 390000), datetime(2018, 7, 25, 17, 15, 6, 377000)),
    (2, datetime(2018, 7, 25, 11, 12, 49, 317000), datetime(2018, 7, 25, 11, 12, 48, 883000))
]

df = spark.createDataFrame(data, ['ID', 'max_ts','min_ts'])
df.printSchema()

# root
#  |-- ID: long (nullable = true)
#  |-- max_ts: timestamp (nullable = true)
#  |-- min_ts: timestamp (nullable = true)
Run Code Online (Sandbox Code Playgroud)

您可以通过将时间戳类型列转换为某种double类型来获取以秒为单位的时间,或者通过将该结果乘以 1000 来获取以毫秒为单位的时间(long如果需要整数,则可以选择转换为 to )。例如

df.select(
    F.col('max_ts').cast('double').alias('time_in_seconds'),
    (F.col('max_ts').cast('double') * 1000).cast('long').alias('time_in_milliseconds'),
).toPandas()

#     time_in_seconds  time_in_milliseconds
# 0    1532538906.390         1532538906390
# 1    1532517169.317         1532517169317
Run Code Online (Sandbox Code Playgroud)

最后,如果您想要两次时间之间的差异(以毫秒为单位),您可以这样做:

df.select(
    ((F.col('max_ts').cast('double') - F.col('min_ts').cast('double')) * 1000).cast('long').alias('diff_in_milliseconds'),
).toPandas()

#    diff_in_milliseconds
# 0                    13
# 1                   434
Run Code Online (Sandbox Code Playgroud)

我正在 PySpark 2.4.2 上执行此操作。无需使用任何字符串连接。


Tan*_*jin 6

这是预期的行为- 它在源代码文档字符串unix_timestamp中明确指出它只返回秒,因此在进行计算时会删除毫秒部分。

如果您想进行该计算,可以使用该substring函数连接数字,然后进行差分。请参阅下面的示例。请注意,这假设数据完整,例如毫秒完全满足(所有 3 位数字):

import pyspark.sql.functions as F

timeFmt = "yyyy-MM-dd' 'HH:mm:ss.SSS"
data = [
    (1, '2018-07-25 17:15:06.390', '2018-07-25 17:15:06.377'),  # note the '390'
    (2, '2018-07-25 11:12:49.317', '2018-07-25 11:12:48.883')
]

df = spark.createDataFrame(data, ['ID', 'max_ts', 'min_ts'])\
    .withColumn('max_milli', F.unix_timestamp('max_ts', format=timeFmt) + F.substring('max_ts', -3, 3).cast('float')/1000)\
    .withColumn('min_milli', F.unix_timestamp('min_ts', format=timeFmt) + F.substring('min_ts', -3, 3).cast('float')/1000)\
    .withColumn('diff', (F.col('max_milli') - F.col('min_milli')).cast('float') * 1000)

df.show(truncate=False)

+---+-----------------------+-----------------------+----------------+----------------+---------+
|ID |max_ts                 |min_ts                 |max_milli       |min_milli       |diff     |
+---+-----------------------+-----------------------+----------------+----------------+---------+
|1  |2018-07-25 17:15:06.390|2018-07-25 17:15:06.377|1.53255330639E9 |1.532553306377E9|13.000011|
|2  |2018-07-25 11:12:49.317|2018-07-25 11:12:48.883|1.532531569317E9|1.532531568883E9|434.0    |
+---+-----------------------+-----------------------+----------------+----------------+---------+
Run Code Online (Sandbox Code Playgroud)