# DolphinDB教程:K线计算 DolphinDB提供了功能强大的内存计算引擎,内置时间序列函数,分布式计算以及流数据处理引擎,在众多场景下均可高效的计算K线。本教程将介绍DolphinDB如何通过批量处理和流式处理计算K线。 - 计算历史数据K线 可以指定K线窗口的起始时间;一天中可以存在多个交易时段,包括隔夜时段;K线窗口可重叠;使用交易量作为划分K线窗口的维度。需要读取的数据量特别大并且需要将结果写入数据库时,可使用DolphinDB内置的Map-Reduce函数并行计算。 - 实时计算K线 使用API实时接收市场数据,并使用DolphinDB内置的流数据时序计算引擎(time-series aggregator)进行实时计算得到K线数据。 ## 1. 历史数据K线计算 使用历史数据计算K线,可使用DolphinDB的内置函数[`bar`](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/b/bar.html),[`dailyAlignedBar`](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/d/dailyAlignedBar.html),或[`wj`](https://www.dolphindb.cn/cn/help/SQLStatements/TableJoiners/windowjoin.html)。 ### 1.1 不指定K线窗口的起始时刻 这种情况可使用`bar`函数。bar(X,Y)返回X减去X除以Y的余数(X-mod(X,Y)),一般用于将数据分组。如下例所示。 ``` date = 09:32m 09:33m 09:45m 09:49m 09:56m 09:56m; bar(date, 5); ``` 返回结果: ``` [09:30m,09:30m,09:45m,09:45m,09:55m,09:55m] ``` **例子1**:使用以下数据模拟美国股票市场: ``` n = 1000000 date = take(2019.11.07 2019.11.08, n) time = (09:30:00.000 + rand(int(6.5*60*60*1000), n)).sort!() timestamp = concatDateTime(date, time) price = 100+cumsum(rand(0.02, n)-0.01) volume = rand(1000, n) symbol = rand(`AAPL`FB`AMZN`MSFT, n) trade = table(symbol, date, time, timestamp, price, volume).sortBy!(`symbol`timestamp) undef(`date`time`timestamp`price`volume`symbol) ``` 计算5分钟K线: ``` barMinutes = 5 OHLC = select first(price) as open, max(price) as high, min(price) as low, last(price) as close, sum(volume) as volume from trade group by symbol, date, bar(time, barMinutes*60*1000) as barStart ``` 请注意,以上数据中,time列的精度为毫秒。若time列精度不是毫秒,则应当将 barMinutes\*60*1000 中的数字做相应调整。 ### 1.2 指定K线窗口的起始时刻 需要指定K线窗口的起始时刻,可使用`dailyAlignedBar`函数。该函数可处理每日多个交易时段,亦可处理隔夜时段。 请注意,使用`dailyAlignedBar`函数时,时间列必须含有日期信息,包括 DATETIME, TIMESTAMP 或 NANOTIMESTAMP 这三种类型的数据。指定每个交易时段起始时刻的参数 timeOffset 必须使用相应的去除日期信息之后的 SECOND,TIME 或 NANOTIME 类型的数据。 **例子2**(每日一个交易时段):计算美国股票市场7分钟K线。数据沿用例子1中的trade表。 ``` barMinutes = 7 OHLC = select first(price) as open, max(price) as high, min(price) as low, last(price) as close, sum(volume) as volume from trade group by symbol, dailyAlignedBar(timestamp, 09:30:00.000, barMinutes*60*1000) as barStart ``` **例子3**(每日两个交易时段):中国股票市场每日有两个交易时段,上午时段为9:30至11:30,下午时段为13:00至15:00。 使用以下脚本产生模拟数据: ``` n = 1000000 date = take(2019.11.07 2019.11.08, n) time = (09:30:00.000 + rand(2*60*60*1000, n/2)).sort!() join (13:00:00.000 + rand(2*60*60*1000, n/2)).sort!() timestamp = concatDateTime(date, time) price = 100+cumsum(rand(0.02, n)-0.01) volume = rand(1000, n) symbol = rand(`600519`000001`600000`601766, n) trade = table(symbol, timestamp, price, volume).sortBy!(`symbol`timestamp) undef(`date`time`timestamp`price`volume`symbol) ``` 计算7分钟K线: ``` barMinutes = 7 sessionsStart=09:30:00.000 13:00:00.000 OHLC = select first(price) as open, max(price) as high, min(price) as low, last(price) as close, sum(volume) as volume from trade group by symbol, dailyAlignedBar(timestamp, sessionsStart, barMinutes*60*1000) as barStart ``` **例子4**(每日两个交易时段,包含隔夜时段):某些期货每日有多个交易时段,且包括隔夜时段。本例中,第一个交易时段为上午8:45到下午13:45,另一个时段为隔夜时段,从下午15:00到第二天上午05:00。 使用以下脚本产生模拟数据: ``` daySession = 08:45:00.000 : 13:45:00.000 nightSession = 15:00:00.000 : 05:00:00.000 n = 1000000 timestamp = rand(concatDateTime(2019.11.06, daySession[0]) .. concatDateTime(2019.11.08, nightSession[1]), n).sort!() price = 100+cumsum(rand(0.02, n)-0.01) volume = rand(1000, n) symbol = rand(`A120001`A120002`A120003`A120004, n) trade = select * from table(symbol, timestamp, price, volume) where timestamp.time() between daySession or timestamp.time()>=nightSession[0] or timestamp.time(), `symbol`time) ``` ### 1.4 使用交易量划分K线窗口 上面的例子我们均使用时间作为划分K线窗口的维度。在实践中,也可以使用其他变量,譬如用累计的交易量作为划分K线窗口的依据。 **例子6** (使用累计的交易量计算K线):交易量每增加1000000计算一次K线。 ``` n = 1000000 sampleDate = 2019.11.07 symbols = `600519`000001`600000`601766 trade = table(take(sampleDate, n) as date, (09:30:00.000 + rand(7200000, n/2)).sort!() join (13:00:00.000 + rand(7200000, n/2)).sort!() as time, rand(symbols, n) as symbol, 100+cumsum(rand(0.02, n)-0.01) as price, rand(1000, n) as volume) volThreshold = 1000000 t = select first(time) as barStart, first(price) as open, max(price) as high, min(price) as low, last(price) as close, last(cumvol) as cumvol from (select symbol, time, price, cumsum(volume) as cumvol from trade context by symbol) group by symbol, bar(cumvol, volThreshold) as volBar ``` 代码采用了嵌套查询的方法。子查询为每个股票生成累计的交易量cumvol,然后在主查询中根据累计的交易量用`bar`函数生成窗口。 ### 1.5 使用MapReduce函数加速 若需从数据库中提取较大量级的历史数据,计算K线,然后存入数据库,可使用DolphinDB内置的Map-Reduce函数[`mr`](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/m/mr.html)进行数据的并行读取与计算。这种方法可以显著提高速度。 本例使用美国股票市场的交易数据。原始数据存于"dfs://TAQ"数据库的"trades"表中。"dfs://TAQ"数据库采用复合分区:基于交易日期Date的值分区与基于股票代码Symbol的范围分区。 (1) 将存于磁盘的原始数据表的元数据载入内存: ``` login(`admin, `123456) db = database("dfs://TAQ") trades = db.loadTable("trades") ``` (2) 在磁盘上创建一个空的数据表,以存放计算结果。以下代码建立一个模板表(model),并根据此模板表的schema在数据库"dfs://TAQ"中创建一个空的 OHLC 表以存放K线计算结果: ``` model=select top 1 Symbol, Date, Time.second() as bar, PRICE as open, PRICE as high, PRICE as low, PRICE as close, SIZE as volume from trades where Date=2007.08.01, Symbol=`EBAY if(existsTable("dfs://TAQ", "OHLC")) db.dropTable("OHLC") db.createPartitionedTable(model, `OHLC, `Date`Symbol) ``` (3) 使用`mr`函数计算K线数据,并将结果写入 OHLC 表中: ``` def calcOHLC(inputTable){ tmp=select first(PRICE) as open, max(PRICE) as high, min(PRICE) as low, last(PRICE) as close, sum(SIZE) as volume from inputTable where Time.second() between 09:30:00 : 15:59:59 group by Symbol, Date, 09:30:00+bar(Time.second()-09:30:00, 5*60) as bar loadTable("dfs://TAQ", `OHLC).append!(tmp) return tmp.size() } ds = sqlDS(