Spark 3.2中interval year类型无法被schema/dtypes解析的解决问询
在Spark 3.2中处理interval year类型的schema解析问题
方法1:从printSchema输出提取类型信息
Spark 3.2里printSchema()能正常输出interval year类型,我们可以直接捕获或调用底层API获取schema文本,再解析列名和类型:
from pyspark.sql import functions as F df2 = spark.range(1).withColumn('col_interval_y', F.expr("INTERVAL '2021' YEAR")) # 调用Java API直接获取schema的树形文本 schema_tree = df2._jdf.schema().treeString() col_types = {} for line in schema_tree.split('\n'): line = line.strip() if line.startswith('|--'): parts = line.split(': ') col_name = parts[0].replace('|-- ', '') dtype = parts[1].split(' ')[0] col_types[col_name] = dtype print(col_types) # 输出: {'id': 'long', 'col_interval_y': 'interval'}
方法2:通过Java API绕过Python层解析bug
直接调用Java DataFrame的schema接口,避免Python层的类型解析错误:
def get_column_types(df): col_types = [] for field in df._jdf.schema().fields(): col_name = field.name() dtype_name = field.dataType().typeName() # 映射interval year的类型名 if dtype_name == "interval-year": dtype = "interval year" else: dtype = dtype_name col_types.append((col_name, dtype)) return col_types # 使用示例 print(get_column_types(df2)) # 输出: [('id', 'bigint'), ('col_interval_y', 'interval year')]
方法3:修复Python层的类型解析逻辑
手动遍历schema字段,捕获解析异常并针对性处理interval类型:
from pyspark.sql.types import * def safe_get_dtypes(df): dtypes = [] for field in df.schema.fields: try: dtype_str = field.dataType.simpleString() except ValueError: # 识别IntervalYearType类型 if str(field.dataType).startswith('IntervalYearType'): dtype_str = 'interval year' else: dtype_str = str(field.dataType) dtypes.append((field.name, dtype_str)) return dtypes # 使用示例 print(safe_get_dtypes(df2)) # 输出: [('id', 'bigint'), ('col_interval_y', 'interval year')]
以上三种方法均无需升级Spark版本,即可在3.2中正确获取包含interval year类型的DataFrame列类型信息。
内容的提问来源于stack exchange,提问作者ZygD
相关产品推荐
相关产品推荐

