You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PyIceberg分区演化写入Minio-Iceberg时遇ValueError问题求助

PyIceberg写入Minio-Iceberg分区演化错误排查与解决

错误信息

ValueError: Could not find in old schema: 1000: datetime_day: day(15)

问题代码

schema = Schema(
        NestedField(field_id=1, name="id", field_type=IntegerType(), required=True),
        NestedField(field_id=2, name="uid", field_type=StringType(), required=True),
        NestedField(field_id=3, name="password", field_type=StringType(), required=True),
        NestedField(field_id=4, name="first_name", field_type=StringType(), required=True),
        NestedField(field_id=5, name="last_name", field_type=StringType(), required=True),
        NestedField(field_id=6, name="username", field_type=StringType(), required=True),
        NestedField(field_id=7, name="email", field_type=StringType(), required=True),
        NestedField(field_id=8, name="avatar", field_type=StringType(), required=True),
        NestedField(field_id=9, name="gender", field_type=StringType(), required=True),
        NestedField(field_id=10, name="phone_number", field_type=StringType(), required=True),
        NestedField(field_id=11, name="social_insurance_number", field_type=StringType(), required=True),
        NestedField(field_id=12, name="date_of_birth", field_type=TimestampType(), required=True),
        NestedField(field_id=13, name="credit_card_cc_number", field_type=StringType(), required=True),
        NestedField(field_id=14, name="row_id", field_type=StringType(), required=True),
        NestedField(field_id=15, name="datetime_day", field_type=StringType(), required=True)
    )

partition_spec = PartitionSpec(
    PartitionField(
        source_id=15, field_id=1000, transform=DayTransform(), name="datetime_day"
    )
)

sort_order = SortOrder(SortField(source_id=1, transform=IdentityTransform()))

hive_catalog = load_catalog(
    "hive",  
    **{
        "uri": "thrift://hive-metastore:9083",
        "s3.endpoint": "http://minio:9000",
        "s3.access-key-id": "my-ak",
        "s3.secret-access-key": "my-sak"
    }
)

arrow_df = pa.Table.from_pandas(wh_table['dim_user_info'])
if not hive_catalog.table_exists("lakehouse_dw1.dim_user_info"):
    table = hive_catalog.create_table(
        "lakehouse_dw1.dim_user_info",
        schema=arrow_df.schema,
        partition_spec=partition_spec,
        sort_order=sort_order
    )
else:
    table = hive_catalog.load_table("lakehouse_dw1.dim_user_info")
    table.append(arrow_df)

错误原因

  1. Schema不一致:创建表时传入arrow_df.schema而非自定义的schema,导致自定义schema中定义的字段ID(包括分区字段的field_id=1000)未被注册到表元数据中。后续追加数据时,分区规范中的field_id=1000在已存在的表schema中找不到,触发错误。
  2. 分区转换逻辑错误:自定义schema中datetime_day是StringType,但分区使用DayTransform——该转换要求源字段必须是TimestampType或DateType,当前配置会导致分区转换失败,属于潜在问题。

解决办法

1. 统一表创建的Schema来源

创建表时使用自定义的schema,确保字段ID和类型被正确记录到表元数据:

if not hive_catalog.table_exists("lakehouse_dw1.dim_user_info"):
    table = hive_catalog.create_table(
        "lakehouse_dw1.dim_user_info",
        schema=schema,  # 替换为自定义Schema
        partition_spec=partition_spec,
        sort_order=sort_order
    )

2. 修正分区转换的源字段

将分区源字段改为日期/时间类型的字段,比如已有的date_of_birth(TimestampType):

partition_spec = PartitionSpec(
    PartitionField(
        source_id=12,  # 对应date_of_birth的字段ID
        field_id=1000,
        transform=DayTransform(),
        name="date_of_birth_day"
    )
)

如果必须使用datetime_day,需将其类型改为DateType或TimestampType,并同步调整DataFrame中的对应字段类型。

3. 清理错误创建的表

如果表已经错误创建,需先删除旧表再重新创建:

if hive_catalog.table_exists("lakehouse_dw1.dim_user_info"):
    hive_catalog.drop_table("lakehouse_dw1.dim_user_info")
# 执行上述修正后的创建表逻辑

内容的提问来源于stack exchange,提问作者Phương Vũ duy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 10:52:28