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)
错误原因
- Schema不一致:创建表时传入
arrow_df.schema而非自定义的schema,导致自定义schema中定义的字段ID(包括分区字段的field_id=1000)未被注册到表元数据中。后续追加数据时,分区规范中的field_id=1000在已存在的表schema中找不到,触发错误。 - 分区转换逻辑错误:自定义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
相关产品推荐
相关产品推荐

