Python双路径写入DolphinDB键值流表时主键未去重问题排查
DolphinDB键值流表写入重复行问题排查与解决
问题背景
从Python向DolphinDB 2.00.16.1版本的键值流表写入实时K线数据,以symbol和unixTime作为复合主键,但查询时始终出现重复行。尝试切换为enableTableShareAndCachePurge、将复合主键改为date+time,问题依旧。
表结构定义
colNames = `symbol`exchange`tradingDay`date`time`open`high`low`close`volume`turnover`unixTime colTypes = [SYMBOL,SYMBOL,DATE,DATE,TIME,DOUBLE,DOUBLE,DOUBLE,DOUBLE,LONG,DOUBLE,LONG] enableTableShareAndPersistence( keyedStreamTable(`symbol`unixTime, 100000:0, colNames, colTypes), `stock_candle_stream, cacheSize=80000 )
数据写入路径
数据通过两条独立路径写入:
- 批量写入:一次性推送完整DataFrame
self.session.run("tableInsert{objByName('" + stream_table + "')}", data)
- 逐行写入:单独插入每一根K线
self.writer.insert( bar.symbol, mapp.get(bar.exchange, ''), day, bar.datetime.date(), _time, round(bar.open_price, 2), round(bar.high_price, 2), round(bar.low_price, 2), round(bar.close_price, 2), int(bar.volume), bar.amount, int(bar.datetime.timestamp()) )
核心原因:时间戳精度不匹配
两条写入路径的unixTime精度不一致:
- 批量写入的
unixTime为毫秒级(13位数字),pandas会保留毫秒精度 - 逐行写入时,
int(bar.datetime.timestamp())将时间戳转换为秒级(10位数字),导致同一K线的主键值相差1000倍,键值流表的去重机制无法触发。
复现示例
DolphinDB端复现脚本
// ===================================================== // 复现示例:键值流表去重失效 // ===================================================== // 1. 清理已有表 undef(`stock_candle_stream, SHARED) // 2. 创建相同表结构 colNames = `symbol`exchange`tradingDay`date`time`open`high`low`close`volume`turnover`unixTime colTypes = [SYMBOL,SYMBOL,DATE,DATE,TIME,DOUBLE,DOUBLE,DOUBLE,DOUBLE,LONG,DOUBLE,LONG] share keyedStreamTable(`symbol`unixTime, 100000:0, colNames, colTypes) as stock_candle_stream // 3. 设置随机种子保证结果可复现 setRandomSeed(42) // 4. 生成5行毫秒级时间戳数据(模拟批量写入路径) unixTime_base_ms = 1746000000000L batchData = table( take(`AAPL, 5) as symbol, take(`NASDAQ, 5) as exchange, take(2026.05.18, 5) as tradingDay, take(2026.05.18, 5) as date, take(10:00:00, 5) + 0..4 as time, rand(100.0, 5) as open, rand(100.0, 5) as high, rand(100.0, 5) as low, rand(100.0, 5) as close, rand(10000, 5) as volume, rand(1000000.0, 5) as turnover, unixTime_base_ms + 0..4 * 60000 as unixTime ) // 插入批量数据——5行,13位时间戳 stock_candle_stream.append!(batchData) // 5. 生成相同K线但秒级时间戳数据(模拟逐行写入路径) unixTime_base_s = 1746000000L rowData = table( take(`AAPL, 5) as symbol, take(`NASDAQ, 5) as exchange, take(2026.05.18, 5) as tradingDay, take(2026.05.18, 5) as date, take(10:00:00, 5) + 0..4 as time, rand(100.0, 5) as open, rand(100.0, 5) as high, rand(100.0, 5) as low, rand(100.0, 5) as close, rand(10000, 5) as volume, rand(1000000.0, 5) as turnover, unixTime_base_s + 0..4 * 60 as unixTime ) // 插入逐行数据——5行,10位时间戳 stock_candle_stream.append!(rowData) // 6. 统计行数——结果为10行而非5行 // 尽管每组行代表同一K线,但时间戳值相差1000倍,未触发去重 select count(*) from stock_candle_stream // 预期结果:10行(5行批量+5行逐行,无去重) // 7. 查看重复行——对比时间戳值 select symbol, time, open, high, low, close, unixTime, strFormat(unixTime) as unixTime_raw from stock_candle_stream where symbol = `AAPL order by time, unixTime // 注意:同一K线的unixTime显示为1746000000(秒级)和1746000000000(毫秒级) // ===================================================== // 验证:精度一致时去重正常 // ===================================================== // 插入两次相同时间戳的K线——去重生效 unixTime_test = 1746000060000L insert into stock_candle_stream values(`AAPL, `NASDAQ, 2026.05.18, 2026.05.18, 10:01:00, 101.0, 102.0, 99.5, 101.5, 8500, 862750.0, unixTime_test) insert into stock_candle_stream values(`AAPL, `NASDAQ, 2026.05.18, 2026.05.18, 10:01:00, 101.0, 102.0, 99.5, 101.5, 8500, 862750.0, unixTime_test) // 仅插入1行(第二次插入被自动丢弃) select symbol, time, unixTime from stock_candle_stream where unixTime = 1746000060000L // 预期结果:1行
Python端复现脚本
import numpy as np import dolphindb as ddb s = ddb.session() s.connect("localhost", 8848) # 先运行DolphinDB端的表创建脚本,再执行以下代码 # 路径A——批量写入毫秒级时间戳(13位) np.random.seed(42) n = 5 unixTime_ms = [1746000000000 + i * 60000 for i in range(n)] s.run(""" data = table( take(`AAPL, {n}) as symbol, take(`NASDAQ, {n}) as exchange, take(2026.05.18, {n}) as tradingDay, take(2026.05.18, {n}) as date, take(10:00:00, {n}) + 0..{e} as time, {open} as open, {high} as high, {low} as low, {close} as close, {volume} as volume, {turnover} as turnover, {unixTime} as unixTime ) """.format( n=n, e=n-1, open=list(np.random.uniform(100, 101, n).round(2)), high=list(np.random.uniform(101, 102, n).round(2)), low=list(np.random.uniform(99, 100, n).round(2)), close=list(np.random.uniform(100, 101, n).round(2)), volume=list(np.random.randint(5000, 15000, n)), turnover=list(np.random.uniform(500000, 1500000, n).round(2)), unixTime=unixTime_ms, )) s.run("stock_candle_stream.append!(data)") # 路径B——逐行写入秒级时间戳(10位) unixTime_s = [1746000000 + i * 60 for i in range(n)] for i in range(n): s.run(""" insert into stock_candle_stream values(`AAPL, `NASDAQ, 2026.05.18, 2026.05.18, 10:00:00, {o}, {h}, {l}, {c}, {v}, {t}, {u}) """.format( o=np.round(np.random.uniform(100, 101), 2), h=np.round(np.random.uniform(101, 102), 2), l=np.round(np.random.uniform(99, 100), 2), c=np.round(np.random.uniform(100, 101), 2), v=int(np.random.randint(5000, 15000)), t=np.round(np.random.uniform(500000, 1500000), 2), u=unixTime_s[i], )) # 检查结果——应返回10行而非5行 print(s.run("select count(*) from stock_candle_stream"))
解决方案
统一两条写入路径的时间戳精度,确保unixTime值完全一致:
- 方案1:将逐行写入的时间戳改为毫秒级,修改Python代码中的时间戳转换逻辑:
int(bar.datetime.timestamp() * 1000) # 转换为毫秒级时间戳 - 方案2:将批量写入的时间戳转换为秒级,确保与逐行写入保持一致。
内容的提问来源于stack exchange,提问作者Kent
相关产品推荐
相关产品推荐

