PySpark日期动态转换异常:Date2/Date3字段返回null求助
问题:动态转换字符串日期列时Date2、Date3返回Null
尝试将字符串形式的各类日期列动态转换为日期类型,代码能识别日期格式,但Date2和Date3字段始终返回null值,希望得到修正方法。
原代码如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, min, max from pyspark.sql.types import IntegerType, FloatType, TimestampType, StringType, DateType from datetime import datetime, date # Define the function to convert values def convert_value(value): try: return int(value) except ValueError: pass try: return float(value) except ValueError: pass datetime_formats = [ '%m/%d/%Y %H:%M:%S', '%Y-%m-%d %H:%M:%S', '%Y-%m-%dT%H:%M:%S', '%Y-%m-%d %H:%M:%S.%f', '%Y-%m-%dT%H:%M:%S.%f' ] for fmt in datetime_formats: try: return datetime.strptime(value, fmt) except ValueError: pass date_formats = [ '%Y-%m-%d', '%d-%m-%Y', '%m/%d/%Y', '%d/%m/%Y', '%Y/%m/%d', '%b %d, %Y', '%d %b %Y' ] for fmt in date_formats: try: return datetime.strptime(value, fmt).date() except ValueError: pass return value # Function to infer data type for each column def infer_column_type(df, column): min_value = df.select(min(col(column))).collect()[0][0] max_value = df.select(max(col(column))).collect()[0][0] for value in [min_value, max_value]: if value is not None: converted_value = convert_value(value) print(f"Column: {column}, Value: {value}, Converted: {converted_value}") # Debug print if isinstance(converted_value, int): return IntegerType() elif isinstance(converted_value, float): return FloatType() elif isinstance(converted_value, datetime): return TimestampType() elif isinstance(converted_value, date): return DateType() return StringType() # Example data with different date formats in separate columns data = [ ('1', '2021-01-01', '01-02-2021', '1/2/2021', '2021-01-01T12:34:56', '1.1', 1), ('2', '2021-02-01', '02-03-2021', '2/3/2021', '2021-02-01T13:45:56', '2.2', 2), ('3', '2021-03-01', '03-04-2021', '3/4/2021', '2021-03-01T14:56:56', '3.3', 3) ] # Create DataFrame spark = SparkSession.builder.appName("example").getOrCreate() columns = ['A', 'Date1', 'Date2', 'Date3', 'Date4', 'C', 'D'] df = spark.createDataFrame(data, columns) # Apply inferred data types to columns for column in df.columns: inferred_type = infer_column_type(df, column) df = df.withColumn(column, df[column].cast(inferred_type)) # Show the result df.show() df.dtypes
原因分析
核心问题是自定义函数convert_value仅用于识别列类型,但后续直接调用Spark的cast方法转换列:
- Spark的
cast(DateType())只支持有限的默认日期格式(如yyyy-MM-dd),像01-02-2021(dd-MM-yyyy)、1/2/2021(MM/dd/yyyy)这类非标准格式无法被自动识别,转换失败后返回null。 infer_column_type函数仅推断了列应该是DateType,但没有将自定义的格式匹配逻辑应用到整列的转换过程中。
修正方案
方案1:用自定义UDF复用现有格式匹配逻辑
直接将convert_value封装为UDF,对整列应用格式转换,完全兼容你已有的格式列表:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf from pyspark.sql.types import IntegerType, FloatType, TimestampType, StringType, DateType from datetime import datetime, date def convert_value(value): try: return int(value) except ValueError: pass try: return float(value) except ValueError: pass datetime_formats = [ '%m/%d/%Y %H:%M:%S', '%Y-%m-%d %H:%M:%S', '%Y-%m-%dT%H:%M:%S', '%Y-%m-%d %H:%M:%S.%f', '%Y-%m-%dT%H:%M:%S.%f' ] for fmt in datetime_formats: try: return datetime.strptime(value, fmt) except ValueError: pass date_formats = [ '%Y-%m-%d', '%d-%m-%Y', '%m/%d/%Y', '%d/%m/%Y', '%Y/%m/%d', '%b %d, %Y', '%d %b %Y' ] for fmt in date_formats: try: return datetime.strptime(value, fmt).date() except ValueError: pass return value # 定义UDF,自动推断返回类型 convert_udf = udf(convert_value) data = [ ('1', '2021-01-01', '01-02-2021', '1/2/2021', '2021-01-01T12:34:56', '1.1', 1), ('2', '2021-02-01', '02-03-2021', '2/3/2021', '2021-02-01T13:45:56', '2.2', 2), ('3', '2021-03-01', '03-04-2021', '3/4/2021', '2021-03-01T14:56:56', '3.3', 3) ] spark = SparkSession.builder.appName("example").getOrCreate() columns = ['A', 'Date1', 'Date2', 'Date3', 'Date4', 'C', 'D'] df = spark.createDataFrame(data, columns) # 对所有列应用UDF转换 for column in df.columns: df = df.withColumn(column, convert_udf(col(column))) df.show() print(df.dtypes)
方案2:结合Spark原生函数指定格式转换(性能更优)
改进类型推断逻辑,同时返回列类型和对应的日期格式,再用Spark的to_date/to_timestamp指定格式转换:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, min, max, to_date, to_timestamp from pyspark.sql.types import IntegerType, FloatType, TimestampType, StringType, DateType from datetime import datetime, date # 推断单个值的日期格式和对应类型 def infer_date_format(value): datetime_formats = [ ('%m/%d/%Y %H:%M:%S', TimestampType()), ('%Y-%m-%d %H:%M:%S', TimestampType()), ('%Y-%m-%dT%H:%M:%S', TimestampType()), ('%Y-%m-%d %H:%M:%S.%f', TimestampType()), ('%Y-%m-%dT%H:%M:%S.%f', TimestampType()) ] for fmt, dtype in datetime_formats: try: datetime.strptime(value, fmt) return (fmt, dtype) except ValueError: pass date_formats = [ ('%Y-%m-%d', DateType()), ('%d-%m-%Y', DateType()), ('%m/%d/%Y', DateType()), ('%d/%m/%Y', DateType()), ('%Y/%m/%d', DateType()), ('%b %d, %Y', DateType()), ('%d %b %Y', DateType()) ] for fmt, dtype in date_formats: try: datetime.strptime(value, fmt).date() return (fmt, dtype) except ValueError: pass return (None, StringType()) # 推断列的类型和对应的格式 def infer_column_info(df, column): min_value = df.select(min(col(column))).collect()[0][0] max_value = df.select(max(col(column))).collect()[0][0] for value in [min_value, max_value]: if value is not None: # 先尝试数字类型 try: int(value) return (IntegerType(), None) except ValueError: pass try: float(value) return (FloatType(), None) except ValueError: pass # 尝试日期时间类型 fmt, dtype = infer_date_format(value) if fmt is not None: return (dtype, fmt) return (StringType(), None) data = [ ('1', '2021-01-01', '01-02-2021', '1/2/2021', '2021-01-01T12:34:56', '1.1', 1), ('2', '2021-02-01', '02-03-2021', '2/3/2021', '2021-02-01T13:45:56', '2.2', 2), ('3', '2021-03-01', '03-04-2021', '3/4/2021', '2021-03-01T14:56:56', '3.3', 3) ] spark = SparkSession.builder.appName("example").getOrCreate() columns = ['A', 'Date1', 'Date2', 'Date3', 'Date4', 'C', 'D'] df = spark.createDataFrame(data, columns) # 应用转换逻辑 for column in df.columns: dtype, fmt = infer_column_info(df, column) if dtype in [DateType(), TimestampType()] and fmt is not None: if dtype == DateType(): df = df.withColumn(column, to_date(col(column), fmt)) else: df = df.withColumn(column, to_timestamp(col(column), fmt)) else: df = df.withColumn(column, col(column).cast(dtype)) df.show() print(df.dtypes)
内容的提问来源于stack exchange,提问作者user3486773
相关产品推荐
相关产品推荐

