在AWS Glue中用delta.io创建DeltaTable遇类型错误的解决咨询
问题解决方法
错误原因分析
- Schema定义错误:你把
item_size定义成了DoubleType,但实际该字段是数组类型,应该用ArrayType(DoubleType())。 - 类型不兼容:从Athena读取的pandas DataFrame中,
item_size的元素是numpy.ndarray类型,而Spark的ArrayType(DoubleType)只接受Python原生list,类型不匹配导致报错。
具体解决方案
方案1:修正Schema并转换numpy数组为list
先修正你的Spark Schema定义,再处理pandas DataFrame中的类型问题:
- 修正Schema:
from pyspark.sql.types import StructType, StructField, IntegerType, ArrayType, DoubleType, TimestampType registration_schema = StructType([ StructField("item_id", IntegerType(), nullable=True), # 把DoubleType改为ArrayType(DoubleType()) StructField("item_size", ArrayType(DoubleType(), nullable=True), nullable=True), # 按实际情况补充其他字段 StructField("acquiredate", TimestampType(), nullable=True), StructField("acquiredate_tz", TimestampType(), nullable=True) ])
- 转换numpy数组为Python list:
import numpy as np # 处理item_size列,空值保持不变 df["item_size"] = df["item_size"].apply(lambda x: x.tolist() if isinstance(x, np.ndarray) else x) # 如果box_size也是同样的numpy数组类型,同样处理 df["box_size"] = df["box_size"].apply(lambda x: x.tolist() if isinstance(x, np.ndarray) else x) # 重新创建Spark DataFrame df_delta = spark.createDataFrame(df, schema=registration_schema)
方案2:直接用Spark读取Athena(绕开pandas类型问题)
如果不需要中间用pandas处理数据,可以直接用Spark读取Athena,避免类型转换的麻烦:
df_delta = spark.read.format("awsathena") \ .option("database", "db_name") \ .option("query", query) \ .load() # 之后可以直接将df_delta写入Delta Table # df_delta.write.format("delta").save("s3://your-path/")
内容的提问来源于stack exchange,提问作者Cotrariello
相关产品推荐
相关产品推荐

