peewee无主键或唯一约束场景下实现upsert操作的问题咨询
问题根因
你遇到的报错本质是SQLite的ON CONFLICT语法要求指定的冲突目标必须对应表上已有的主键或者唯一约束,你当前的Candles表没有对「交易对+时间周期+开盘时间」这个K线唯一标识加联合唯一约束,所以SQLite无法识别冲突规则。你之前调用on_conflict_replace、on_conflict_ignore也会重复插入,同样是因为没有唯一约束,SQLite无法判断哪些是需要处理的冲突行。
推荐方案(修改表结构后用原生upsert)
这是性能最高、并发安全的方案,步骤如下:
- 修改
Candles模型,添加联合唯一约束
from peewee import UniqueConstraint # 头部导入对应类 class Candles(BaseModel): timeframe = ForeignKeyField(Timeframes) symbol = ForeignKeyField(Symbols) open_time = DateTimeField() open = FloatField() high = FloatField() low = FloatField() close = FloatField() close_time = DateTimeField() base_volume = IntegerField() quote_volume = IntegerField() class Meta: constraints = [ # 三个字段联合唯一,作为一根K线的唯一识别规则 UniqueConstraint(fields=['timeframe', 'symbol', 'open_time']) ]
如果是已经存在数据的旧表,需要先执行SQL添加约束(执行前请先确认现有数据不存在三个字段重复的记录,否则约束添加失败):
ALTER TABLE candles ADD CONSTRAINT unique_candle UNIQUE (timeframe_id, symbol_id, open_time);
- 简化你的upsert代码,冲突目标仅保留三个唯一字段,删除多余的
open字段和preserve参数:
query = (Candles .insert( timeframe = stream['data']['k']['i'], symbol = stream['data']['k']['s'], open_time = stream['data']['k']['t'], open = stream['data']['k']['o'], high = stream['data']['k']['h'], low = stream['data']['k']['l'], close = stream['data']['k']['c'], close_time = stream['data']['k']['T'], base_volume = stream['data']['k']['v'], quote_volume = stream['data']['k']['q'] ) .on_conflict( conflict_target = [Candles.timeframe, Candles.symbol, Candles.open_time], update = { Candles.high: stream['data']['k']['h'], Candles.low: stream['data']['k']['l'], Candles.close: stream['data']['k']['c'], Candles.close_time: stream['data']['k']['T'], Candles.base_volume: stream['data']['k']['v'], Candles.quote_volume: stream['data']['k']['q'] } ) .execute() )
替代方案(无需修改表结构,手动实现upsert)
如果不想改动表结构,可以用「先查询再判断插入/更新」的逻辑,仅适合单线程写入的场景,并发场景下有概率出现重复插入:
# 先查询对应K线是否存在 existing_candle = Candles.get_or_none( timeframe = stream['data']['k']['i'], symbol = stream['data']['k']['s'], open_time = stream['data']['k']['t'] ) if existing_candle: # 存在则更新字段 existing_candle.high = stream['data']['k']['h'] existing_candle.low = stream['data']['k']['l'] existing_candle.close = stream['data']['k']['c'] existing_candle.close_time = stream['data']['k']['T'] existing_candle.base_volume = stream['data']['k']['v'] existing_candle.quote_volume = stream['data']['k']['q'] existing_candle.save() else: # 不存在则插入 Candles.create( timeframe = stream['data']['k']['i'], symbol = stream['data']['k']['s'], open_time = stream['data']['k']['t'], open = stream['data']['k']['o'], high = stream['data']['k']['h'], low = stream['data']['k']['l'], close = stream['data']['k']['c'], close_time = stream['data']['k']['T'], base_volume = stream['data']['k']['v'], quote_volume = stream['data']['k']['q'] )
内容的提问来源于stack exchange,提问作者Jason
相关产品推荐
相关产品推荐

