使用delta-rs写入Delta表时如何添加新列?
Delta Lake使用delta-rs添加新列的报错解决方法
问题描述
使用delta-rs向Delta Lake写入数据,已成功创建包含timestamp、current、voltage、temperature四列的表。尝试添加新列pressure时,修改了Schema并添加overwrite_schema=True参数,但始终报错。错误信息显示表Schema包含pressure列,但数据Schema没有,与实际情况完全相反——数据明明包含pressure列。
初始写入代码
import time import numpy as np import pandas as pd import pyarrow as pa from deltalake.writer import write_deltalake num_rows = 10 timestamp = np.array([time.time() + i * 0.01 for i in range(num_rows)]) current = np.random.rand(num_rows) * 10 voltage = np.random.rand(num_rows) * 100 temperature = np.random.rand(num_rows) * 50 data = { "timestamp": timestamp, "current": current, "voltage": voltage, "temperature": temperature, } df = pd.DataFrame(data) storage_options = { "AWS_DEFAULT_REGION": "us-west-2", "AWS_ACCESS_KEY_ID": "xxx", "AWS_SECRET_ACCESS_KEY": "xxx", "AWS_S3_ALLOW_UNSAFE_RENAME": "true", } schema = pa.schema( [ ("timestamp", pa.float64()), ("current", pa.float64()), ("voltage", pa.float64()), ("temperature", pa.float64()), ] ) write_deltalake( "s3a://my-bucket/delta-tables/motor", df, mode="append", schema=schema, storage_options=storage_options, )
添加pressure列后的写入代码
import time import numpy as np import pandas as pd import pyarrow as pa from deltalake.writer import write_deltalake num_rows = 10 timestamp = np.array([time.time() + i * 0.01 for i in range(num_rows)]) current = np.random.rand(num_rows) * 10 voltage = np.random.rand(num_rows) * 100 temperature = np.random.rand(num_rows) * 50 pressure = np.random.rand(num_rows) * 1000 data = { "timestamp": timestamp, "current": current, "voltage": voltage, "temperature": temperature, "pressure": pressure, } df = pd.DataFrame(data) storage_options = { "AWS_DEFAULT_REGION": "us-west-2", "AWS_ACCESS_KEY_ID": "xxx", "AWS_SECRET_ACCESS_KEY": "xxx", "AWS_S3_ALLOW_UNSAFE_RENAME": "true", } schema = pa.schema( [ ("timestamp", pa.float64()), ("current", pa.float64()), ("voltage", pa.float64()), ("temperature", pa.float64()), ("pressure", pa.float64()), # 新增此行 ] ) write_deltalake( "s3a://my-bucket/delta-tables/motor", df, mode="append", schema=schema, storage_options=storage_options, overwrite_schema=True, # 加不加这个参数都会报同样的错 )
报错信息
... Traceback (most recent call last): File "python3.11/site-packages/deltalake/writer.py", line 180, in write_deltalake raise ValueError( ValueError: 数据Schema与表Schema不匹配 表Schema: timestamp: double current: double voltage: double temperature: double pressure: double 数据Schema: timestamp: double current: double voltage: double temperature: double
解决方案
1. 移除手动指定的schema参数,让delta-rs自动推断数据Schema
问题根源大概率是手动指定的Schema与DataFrame实际Schema存在隐性不匹配,或delta-rs同时接收schema参数和DataFrame时的校验逻辑冲突。移除手动Schema后,让delta-rs自动从DataFrame推断,同时保留overwrite_schema=True来允许Schema演化:
import time import numpy as np import pandas as pd from deltalake.writer import write_deltalake num_rows = 10 timestamp = np.array([time.time() + i * 0.01 for i in range(num_rows)]) current = np.random.rand(num_rows) * 10 voltage = np.random.rand(num_rows) * 100 temperature = np.random.rand(num_rows) * 50 pressure = np.random.rand(num_rows) * 1000 data = { "timestamp": timestamp, "current": current, "voltage": voltage, "temperature": temperature, "pressure": pressure, } df = pd.DataFrame(data) storage_options = { "AWS_DEFAULT_REGION": "us-west-2", "AWS_ACCESS_KEY_ID": "xxx", "AWS_SECRET_ACCESS_KEY": "xxx", "AWS_S3_ALLOW_UNSAFE_RENAME": "true", } write_deltalake( "s3a://my-bucket/delta-tables/motor", df, mode="append", storage_options=storage_options, overwrite_schema=True, )
2. 确认Delta表的当前元数据状态
如果之前的错误操作已经修改了表的Schema(比如表中已存在pressure列),可以先读取表Schema确认:
from deltalake import DeltaTable dt = DeltaTable("s3a://my-bucket/delta-tables/motor", storage_options=storage_options) print(dt.schema().pretty_print())
若表Schema已包含pressure,只需确保DataFrame包含该列,使用mode="append"即可(无需overwrite_schema,因为Schema已兼容)。
3. 升级delta-rs到最新版本
部分旧版本delta-rs在处理Schema演化时存在bug,执行以下命令升级:
pip install --upgrade deltalake
内容的提问来源于stack exchange,提问作者Hongbo Miao
相关产品推荐
相关产品推荐

