Spark创建模拟DataFrame时硬编码Timestamp值报错,如何解决?
解决Spark中硬编码Timestamp值创建DataFrame的报错问题
你编写的Spark代码在创建带TimestampType字段的DataFrame时触发了类型错误,报错信息显示字符串无法直接被TimestampType接收:
TypeError: field fecha_vencimiento: TimestampType can not accept object '2022-10-31 16:00:00.05' in type <class 'str'>
这是因为Spark在显式指定Schema时,不会自动将字符串转换为Timestamp类型,需要显式处理时间值,以下是三种可行的解决方案:
方法1:使用Python datetime对象替代字符串
导入Python的datetime模块,将时间字符串转换为datetime对象,Spark会自动识别并映射到TimestampType:
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType from pyspark.sql import SparkSession from datetime import datetime sc = SparkSession.builder.appName('Testing').getOrCreate() # 将时间字符串转为datetime对象 data = [ ("FIX", 100.0, 0.01, 0.04, 0.05, datetime.strptime('2022-10-31 16:00:00.05', '%Y-%m-%d %H:%M:%S.%f')), ("FLT", 1000.0, 0.02, 0.03, 0.03, datetime.strptime('2022-09-21 08:59:00.00', '%Y-%m-%d %H:%M:%S.%f')), ("MIX", 10000.0, 0.03, 0.02, 0.01, datetime.strptime('2022-08-15 13:25:16.00', '%Y-%m-%d %H:%M:%S.%f')), ("MIX", 15000.0, 0.04, 0.01, 0.04, datetime.strptime('2022-06-03 19:27:03.00', '%Y-%m-%d %H:%M:%S.%f')) ] schema = StructType([ StructField("type", StringType(), nullable=True), StructField("remaining_amt", DoubleType(), nullable=True), StructField("fix_rt", DoubleType(), nullable=True), StructField("flt_rt", DoubleType(), nullable=True), StructField("spread", DoubleType(), nullable=True), StructField("end_date", TimestampType(), nullable=True) ]) mock_df = sc.createDataFrame(data=data, schema=schema) # 验证结果 mock_df.printSchema() mock_df.show(truncate=False)
方法2:先创建临时DataFrame,再转换字段类型
先不指定Schema创建临时DataFrame,之后使用Spark的to_timestamp函数将字符串字段转换为Timestamp类型:
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType from pyspark.sql import SparkSession from pyspark.sql.functions import to_timestamp sc = SparkSession.builder.appName('Testing').getOrCreate() data = [ ("FIX", 100.0, 0.01, 0.04, 0.05, '2022-10-31 16:00:00.05'), ("FLT", 1000.0, 0.02, 0.03, 0.03, '2022-09-21 08:59:00.00'), ("MIX", 10000.0, 0.03, 0.02, 0.01, '2022-08-15 13:25:16.00'), ("MIX", 15000.0, 0.04, 0.01, 0.04, '2022-06-03 19:27:03.00') ] # 创建临时DataFrame,自动推断字段类型 temp_df = sc.createDataFrame(data, ["type", "remaining_amt", "fix_rt", "flt_rt", "spread", "end_date_str"]) # 转换字符串字段为Timestamp类型 mock_df = temp_df.withColumn("end_date", to_timestamp("end_date_str", "yyyy-MM-dd HH:mm:ss.SS")) \ .drop("end_date_str") # 验证结果 mock_df.printSchema() mock_df.show(truncate=False)
方法3:使用ISO 8601格式的时间字符串
将时间字符串中的空格替换为T,转为ISO 8601格式(如'2022-10-31T16:00:00.05'),部分Spark版本支持自动将该格式字符串映射到TimestampType,但这种方式兼容性不如前两种,推荐优先使用显式转换的方法。
内容的提问来源于stack exchange,提问作者eduardo0
相关产品推荐
相关产品推荐

