# DolphinDB 搭建行情回放服务的最佳实践
一个量化策略在生产(交易)环境中运行时,实时数据的处理通常是由事件驱动的。为确保研发和生产使用同一套代码,通常在研发阶段需将历史数据,严格按照事件发生的时间顺序进行回放,以此模拟交易环境。在 DolphinDB 中,用户通过 `replay` 函数可以实现对静态数据的回放,即将历史数据按照时间顺序以“实时数据”的方式注入流数据表中。对相同时间戳的数据还可以指定额外排序列,使数据回放顺序更接近实时交易场景。
在[历史数据回放](https://gitee.com/dolphindb/Tutorials_CN/blob/master/historical_data_replay.md)、[股票行情回放](https://gitee.com/dolphindb/Tutorials_CN/blob/master/stock_market_replay.md)两篇教程中已经介绍了 DolphinDB 的回放功能,本教程更加侧重于回放功能的工程化实践。本教程将介绍如何基于 DolphinDB 分布式数据库、回放功能以及 DolphinDB API 搭建一个行情数据回放服务,该服务支持多个用户同时通过 C++ 、 Python 等客户端提交数据回放请求。
**目录**
- [1. 基于 DolphinDB 的行情回放服务](#1-基于-dolphindb-的行情回放服务)
- [1.1 行情回放服务架构](#11-行情回放服务架构)
- [1.2 回放服务搭建步骤](#12-回放服务搭建步骤)
- [2. 行情数据分区存储方案](#2-行情数据分区存储方案)
- [3. 行情回放自定义函数](#3-行情回放自定义函数)
- [3.1 stkReplay 函数:行情回放服务主函数](#31-stkreplay-函数行情回放服务主函数)
- [3.2 dsTb 函数:构造回放数据源](#32-dstb-函数构造回放数据源)
- [3.3 replayJob 函数:定义回放任务内容](#33-replayjob-函数定义回放任务内容)
- [3.4 createEnd 函数:构造回放结束信号](#34-createend-函数构造回放结束信号)
- [3.5 封装函数视图](#35-封装函数视图)
- [4. API 提交回放](#4-api-提交回放)
- [4.1 C++ API](#41-c-api)
- [4.2 Python API](#42-python-api)
- [5. API 订阅消费](#5-api-订阅消费)
- [5.1 C++ API](#51-c-api)
- [5.2 Python API](#52-python-api)
- [6. 性能测试](#6-性能测试)
- [6.1 测试服务器配置](#61-测试服务器配置)
- [6.2 50 支深交所股票一天全速并发回放性能测试](#62-50-支深交所股票一天全速并发回放性能测试)
- [6.3 50 支深交所股票跨天全速回放性能测试](#63-50-支深交所股票跨天全速回放性能测试)
- [7. 开发环境配置](#7-开发环境配置)
- [部署 DolphinDB Server](#部署-dolphindb-server)
- [DolphinDB client 开发环境](#dolphindb-client-开发环境)
- [DolphinDB C++ API 安装](#dolphindb-c-api-安装)
- [DolphinDB Python API 安装](#dolphindb-python-api-安装)
- [8. 路线图 (Roadmap)](#8-路线图-roadmap)
- [9. 总结](#9-总结)
- [10. 附录](#10-附录)
# 1. 基于 DolphinDB 的行情回放服务
本教程实现的行情回放服务基于 3 类国内 A 股行情数据源:逐笔委托数据、逐笔成交数据、Level 2 快照数据,支持以下功能与特性:
- C++、Python 客户端提交回放请求(指定回放股票列表 、回放日期、回放速率、回放数据源)
- 多个用户同时回放
- 多个数据源同时有序回放
- 在时间有序的基础上支持排序列有序(如:针对逐笔数据中的交易所原始消息记录号排序)
- 发布回放结束信号
- 对回放结果订阅消费
## 1.1 行情回放服务架构
本教程示例 DolphinDB 搭建的行情回放服务架构如下图所示:
行情回放服务架构
- 行情数据接入:实时行情数据和历史行情数据可以通过 DolphinDB API 或插件存储到 DolphinDB 分布式时序数据库中。
- 函数模块封装:数据查询和回放过程可以通过 DolphinDB 函数视图封装内置,仅暴露股票列表 、回放日期、回放速率、回放数据源等关键参数给行情服务用户。
- 行情用户请求:需要进行行情回放的用户可以通过 DolphinDB 客户端软件(如 DolphinDB GUI 工具、DolphinDB VS Code 插件、DolphinDB API 等)调用封装好的回放函数对存储在数据库中的行情数据进行回放,同时,用户还可以在客户端对回放结果进行实时订阅消费。此外,支持多用户并发回放。
## 1.2 回放服务搭建步骤
本教程示例 DolphinDB 搭建行情回放服务的具体操作步骤如下图所示:
搭建步骤
- Step 1:服务提供者设计合理分区的数据库表,在 DolphinDB 集成开发环境中执行对应的建库建表、数据导入等脚本,以将历史行情数据存储到 DolphinDB 分布式数据库中作为回放服务的数据源。在第二章将给出本教程涉及的 3 类行情数据的分区存储方案及建库建表脚本。
- Step 2:服务提供者在 DolphinDB 集成开发环境中将回放过程中的操作封装成函数视图,通过封装使得行情服务用户不需要关心 DolphinDB 回放功能的细节,只需要指定简单的回放参数(股票、日期、回放速率、数据源)即可提交回放请求。在第三章将给全部函数视图的定义脚本。
- Step 3:服务提供者在外部程序中通过 DolphinDB API 调用上述函数视图实现提交回放的功能。在第四章将给出 API 端提交回放任务的 C++ 实现和 Python 实现。此外,在第四章提交回放的基础上,在第五章将介绍对回放结果的 API 端订阅与消费的代码实现。
在第六章将给出多用户多表回放、多天回放的性能测试结果。最后两章为开发环境配置与总结。
# 2. 行情数据分区存储方案
本教程的回放服务基于 3 类国内 A 股行情数据源:逐笔成交数据、逐笔委托数据、快照数据,均使用 TSDB 存储引擎存储在 DolphinDB 分布式数据库中。
| **数据源** | **代码样例中的分区数据库路径** | **代码样例中的表名** | **分区机制** | **排序列** | **建库建表脚本** |
| :--------- | :----------------------------- | :------------------- | :-------------------------------- | :--------------- | :---------------------------- |
| 逐笔委托 | dfs://Test_order | order | VALUE: 交易日, HASH: [SYMBOL, 25] | 股票ID,交易时间 | 附录 [逐笔委托建库建表脚本](script/appendices_market_replay_bp/order_create.txt) |
| 逐笔成交 | dfs://Test_transaction | transaction | VALUE: 交易日, HASH: [SYMBOL, 25] | 股票ID,交易时间 | 附录 [逐笔成交建库建表脚本](script/appendices_market_replay_bp/transac_create.txt) |
| 快照 | dfs://Test_snapshot | snapshot | VALUE: 交易日, HASH: [SYMBOL, 20] | 股票ID,交易时间 | 附录 [快照建库建表脚本](script/appendices_market_replay_bp/snap_create.txt) |
回放的原理是从数据库中读取需要的行情数据,并根据时间列排序后写入到相应的流数据表。因此,读数据库并排序的性能对回放速度有很大的影响,合理的分区机制将有助于提高数据加载速度。基于用户通常按照日期和股票提交回放请求的特点,设计了上表的分区方案。
此外,本教程的历史数据存储在三节点双副本的 DolphinDB 集群中,集群和副本同样可以提升读取性能,同时可以增加系统的可用性,分区的副本通常是存放在不同的物理节点的,所以一旦某个分区不可用,系统依然可以调用其它副本分区来保证回放服务的正常运转。
附录将提供部分原始数据的 csv 文件(原始行情数据文件)以及对应的示例导入脚本(逐笔委托示例导入脚本、逐笔成交示例导入脚本、快照示例导入脚本),以便读者快速体验搭建本教程所示的行情回放服务。
# 3. 行情回放自定义函数
本章介绍回放过程中的主要函数功能及其实现,最后将函数封装成视图以便通过 API 等方式调用。
本章的开发工具采用 DolphinDB GUI,完整脚本见附录:[行情回放函数](script/appendices_market_replay_bp/replay.txt)。
| 函数名 | 函数入参 | 函数功能 |
| --------- | ------------------------------------------------------------ | ------------------ |
| stkReplay | stkList:回放股票列表
startDate:回放开始日期
endDate:回放结束日期
replayRate:回放速率
replayUuid:回放用户标识
replayName:回放数据源名称 | 行情回放服务主函数 |
| dsTb | timeRS:数据源时间划分
startDate:回放开始日期
endDate:回放结束日期
stkList:回放股票列表
replayName:回放数据源名称 | 构造回放数据源 |
| createEnd | tabName:回放输出表名称
sortColumn:相同时间戳时额外排序列名 | 构造回放结束信号 |
| replayJob | inputDict:回放输入数据源
tabName:回放输出表名称
dateDict:回放排序时间列
timeDict:回放排序时间列
replayRate:回放速率
sortColumn:相同时间戳时额外排序列名 | 定义回放任务内容 |
## 3.1 stkReplay 函数:行情回放服务主函数
函数定义代码:
```
def stkReplay(stkList, mutable startDate, mutable endDate, replayRate, replayUuid, replayName)
{
maxCnt = 50
returnBody = dict(STRING, STRING)
startDate = datetimeParse(startDate, "yyyyMMdd")
endDate = datetimeParse(endDate, "yyyyMMdd") + 1
sortColumn = "ApplSeqNum"
if(stkList.size() > maxCnt)
{
returnBody["errorCode"] = "0"
returnBody["errorMsg"] = "超过单次回放股票上限,最大回放上限:" + string(maxCnt)
return returnBody
}
if(size(replayName) != 0)
{
for(name in replayName)
{
if(not name in ["snapshot", "order", "transaction"])
{
returnBody["errorCode"] = "0"
returnBody["errorMsg"] = "请输入正确的数据源名称,不能识别的数据源名称:" + name
return returnBody
}
}
}
else
{
returnBody["errorCode"] = "0"
returnBody["errorMsg"] = "缺少回放数据源,请输入正确的数据源名称"
return returnBody
}
try
{
if(size(replayName) == 1 && replayName[0] == "snapshot")
{
colName = ["timestamp", "biz_type", "biz_data"]
colType = [TIMESTAMP, SYMBOL, BLOB]
sortColumn = "NULL"
}
else
{
colName = ["timestamp", "biz_type", "biz_data", sortColumn]
colType = [TIMESTAMP, SYMBOL, BLOB, LONG]
}
msgTmp = streamTable(10000000:0, colName, colType)
tabName = "replay_" + replayUuid
enableTableShareAndPersistence(table=msgTmp, tableName=tabName, asynWrite=true, compress=true, cacheSize=10000000, retentionMinutes=60, flushMode=0, preCache=1000000)
timeRS = cutPoints(09:30:00.000..15:00:00.000, 23)
inputDict = dict(replayName, each(dsTb{timeRS, startDate, endDate, stkList}, replayName))
dateDict = dict(replayName, take(`MDDate, replayName.size()))
timeDict = dict(replayName, take(`MDTime, replayName.size()))
jobId = "replay_" + replayUuid
jobDesc = "replay stock data"
submitJob(jobId, jobDesc, replayJob{inputDict, tabName, dateDict, timeDict, replayRate, sortColumn})
returnBody["errorCode"] = "1"
returnBody["errorMsg"] = "后台回放成功"
return returnBody
}
catch(ex)
{
returnBody["errorCode"] = "0"
returnBody["errorMsg"] = "回放行情数据异常,异常信息:" + ex
return returnBody
}
}
```
函数功能:
自定义函数 stkReplay 是整个回放的主体函数,用户传入的参数在 stkReplay 里会进行有效性判断及格式处理,可以根据实际需求更改。
首先,用 maxCnt 来控制用户一次回放股票数量的最大上限,本例中设置的是 50 。returnBody 构造了信息字典,返回给用户以提示执行错误或执行成功。回放开始日期 startDate 和回放结束日期 endDate 利用 [datetimeParse](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/d/datetimeParse.html?highlight=datetimeparse) 函数进行格式处理 。replayRate 是回放速率,replayUuid 是回放表名名称,replayName 是回放数据源列表,sortColumn 是数据源同回放时间戳排序列列名。
当输入参数无误后,便初始化回放结果流表,结果流表为异构流数据表,字段类型为 BLOB 的字段包含了一条原始记录的全部信息,同时结果流表为持久化流表,[enableTableShareAndPersistence](https://www.dolphindb.cn/cn/help/FunctionsandCommands/CommandsReferences/e/enableTableShareAndPersistence.html?highlight=enabletableshareandpersistence) 函数把流数据表共享并把它持久化到磁盘上,使用持久化流表可以避免内存占用过大。当回放数据源包含逐笔成交(transaction )或逐笔委托(order)时,本例实现了对相同时间戳的逐笔数据按交易所原始消息记录号(ApplSeqNum)进行排序(具体实现见 [3.3 replayJob 函数](#33-replayjob-函数定义回放任务内容)),所以结果流表中必须冗余一列来存放排序列。若回放数据源仅包含快照(snapshot)时,则不需要冗余一列排序列。
定义回放需要的其他参数。inputDict 构造了回放数据源列表字典,利用 [each](https://www.dolphindb.cn/cn/help/Functionalprogramming/TemplateFunctions/each.html?highlight=each) 函数和 [部分应用](https://www.dolphindb.cn/cn/help/Functionalprogramming/PartialApplication.html?highlight=部分应用) 可以对多个数据源进行简洁的定义。dateDict 和 timeDict 构造了回放数据源时间戳字典。最后通过 [submitJob](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/s/submitJob.html?highlight=submitjob) 提交后台回放任务。
## 3.2 dsTb 函数:构造回放数据源
函数定义代码:
```
def dsTb(timeRS, startDate, endDate, stkList, replayName)
{
if(replayName == "snapshot"){
tab = loadTable("dfs://Test_snapshot", "snapshot")
}
else if(replayName == "order") {
tab = loadTable("dfs://Test_order", "order")
}
else if(replayName == "transaction") {
tab = loadTable("dfs://Test_transaction", "transaction")
}
else {
return NULL
}
ds = replayDS(sqlObj=