如何通过PyIceberg API向Iceberg表插入新记录?
使用PyIceberg API向Iceberg表插入新记录的实现方法
PyIceberg支持通过append方法向表中追加数据,需将待插入数据转换为PyArrow Table或Pandas DataFrame格式,确保字段名和数据类型与Iceberg表的Schema匹配。以下是针对你的代码的修改实现:
修改后的完整代码
from pyiceberg.catalog import load_catalog import pyarrow as pa import pandas as pd from datetime import datetime import pytz catalog_name = "xxxx" schema_name = "zz" table_name = "yyy" # 加载Catalog和目标表 catalog = load_catalog(name=catalog_name, **{ "uri":"thrift://hive metastore.ip.adress:port", "s3.endpoint": "https://end/point/api/url:port", "py-io-impl": "pyiceberg.io.pyarrow.PyArrowFileIO", "s3.access-key-id": "myaccesskey", "s3.secret-access-key": "secretaccesskey" }) iceberg_table = catalog.load_table(f'{schema_name}.{table_name}') def displayData(): print("Table Schema:\n", iceberg_table.schema()) scan = iceberg_table.scan() print(scan.to_pandas()) def InsertNewRecord(new_record_data): # 将单条/多条记录转换为Pandas DataFrame df = pd.DataFrame(new_record_data if isinstance(new_record_data, list) else [new_record_data]) # 处理时间字段:转换为带UTC时区的datetime类型(匹配Iceberg的TimestampType) df['reading_time'] = pd.to_datetime(df['reading_time'], utc=True) # 调用append方法插入数据 iceberg_table.append(df) print("记录插入成功") # 单条记录插入示例 new_record_data = { "sensor_name": "M2", "reading_time": "2023-10-2 11:48:16+00:00", "value": 20 } # 执行插入并验证 InsertNewRecord(new_record_data) displayData()
关键说明
- 数据格式适配:PyIceberg的
append方法原生支持Pandas DataFrame或PyArrow Table输入,无需手动映射Schema,只要字段名和类型与表结构一致即可。 - 时间类型兼容:Iceberg的
TimestampType通常要求带时区信息,因此必须将字符串时间转换为datetime64[ns, UTC]类型,避免插入时出现类型不兼容错误。 - Append特性:Iceberg的追加操作会生成新的数据文件,不会修改原有数据,保证了数据的不可变性和ACID特性。
批量插入扩展
如果需要插入多条记录,直接传入字典列表即可:
batch_records = [ {"sensor_name": "M2", "reading_time": "2023-10-2 11:48:16+00:00", "value": 20}, {"sensor_name": "M3", "reading_time": "2023-10-2 11:49:16+00:00", "value": 25} ] InsertNewRecord(batch_records)
内容的提问来源于stack exchange,提问作者lion
相关产品推荐
相关产品推荐

