# DolphinDB 通用计算 DolphinDB database不仅可以分布式地存储数据,而且对分布式计算有良好支持。在DolphinDB中,用户可以用系统提供的通用分布式计算框架,通过脚本实现高效的分布式算法,而不需关注具体的底层实现。本文将对DolphinDB通用计算框架中的重要概念和相关函数作出详细解释,并提供丰富的具体使用场景和例子。 - [DolphinDB 通用计算](#dolphindb-通用计算) - [1. 数据源](#1-数据源) - [2. Map-Reduce框架](#2-map-reduce框架) - [2.1. `mr`函数](#21-mr函数) - [2.2. `imr`函数](#22-imr函数) - [3. 数据源相关函数](#3-数据源相关函数) - [3.1. `sqlDS`函数](#31-sqlds函数) - [3.2. `repartitionDS`函数](#32-repartitionds函数) - [3.3. `textChunkDS`函数](#33-textchunkds函数) - [3.4. 第三方数据源提供的数据源接口](#34-第三方数据源提供的数据源接口) - [3.5. 数据源缓存](#35-数据源缓存) - [3.6. 数据源转换](#36-数据源转换) ## 1. 数据源 数据源(data source)是DolphinDB的通用计算框架中的基本概念。它是一种特殊类型的数据对象,是对数据的元描述。通过执行数据源,用户可以获得诸如表、矩阵、向量等数据实体。在DolphinDB的分布式计算框架中,轻量级的数据源对象而不是庞大的数据实体被传输到远程节点,以用于后续的计算,这大大减少了网络流量。 在DolphinDB中,可使用`sqlDS`函数,基于一个SQL表达式产生数据源。这个函数并不直接对表进行查询,而是返回一个或多个SQL子查询的元语句,即数据源。之后,用户可以使用Map-Reduce框架,传入数据源和计算函数,将任务分发到每个数据源对应的节点,并行地完成计算,然后将结果汇总。 关于几种常用的获得数据源的方法,本文的3.1至3.4节中会详细介绍。 ## 2. Map-Reduce框架 Map-Reduce函数是DolphinDB通用分布式计算框架的核心功能。 ### 2.1. `mr`函数 DolphinDB的Map-Reduce函数`mr`的语法是 mr(ds, mapFunc, \[reduceFunc], \[finalFunc], \[parallel=true]),它可接受一组数据源和一个mapFunc函数作为参数。它会将计算任务分发到每个数据源所在的结点,通过mapFunc对每个数据源中的数据进行处理。可选参数reduceFunc会将mapFunc的返回值两两做计算,得到的结果再与第三个mapFunc的返回值计算,如此累积计算,将mapFunc的结果汇总。如果有M个map调用,reduce函数将被调用M-1次。可选参数finalFunc对reduceFunc的返回值做进一步处理。 官方手册中有一个通过`mr`执行分布式最小二乘线性回归的例子。本文通过以下例子,展示如何用一个`mr`调用实现对分布式表每个分区中的数据随机采样十分之一的功能: ``` // 创建数据库和DFS表 db = database("dfs://sampleDB", VALUE, `a`b`c`d) t = db.createPartitionedTable(table(100000:0, `sym`val, [SYMBOL,DOUBLE]), `tb, `sym) n = 3000000 t.append!(table(rand(`a`b`c`d, n) as sym, rand(100.0, n) as val)) // 定义map函数 def sampleMap(t) { sampleRate = 0.1 rowNum = t.rows() sampleIndex = (0..(rowNum - 1)).shuffle()[0:int(rowNum * sampleRate)] return t[sampleIndex] } ds = sqlDS() // 创建数据源 olsEx(ds, `y, `x) // 执行计算 ``` ### 3.2. `repartitionDS`函数 `sqlDS`的数据源是系统自动根据数据的分区而生成的。有时用户需要对数据源做一些限制,例如,在获取数据时,重新指定数据的分区以减少计算量,或者,只需要一部分分区的数据。`repartitionDS`函数就提供了重新划分数据源的功能。 函数`repartitionDS`根据输入的SQL元代码和列名、分区类型、分区方案等,为元代码生成经过重新分区的新数据源。 以下代码提供了一个`repartitionDS`的例子。在这个例子中,DFS表t中有字段deviceId, time, temperature,分别为symbol, datetime和double类型,数据库采用双层分区,第一层对time按VALUE分区,一天一个分区;第二层对deviceId按HASH分成20个区。 现需要按deviceId字段聚合查询95百分位的temperature。如果直接写查询`select percentile(temperature, 95) from t group by deviceId`,由于`percentile`函数没有Map-Reduce实现,这个查询将无法完成。 一个方案是将所需字段全部加载到本地,计算95百分位,但当数据量过大时,计算资源可能不足。`repartitionDS`提供了一个解决方案:将表基于deviceId按其原有分区方案HASH重新分区,每个新的分区对应原始表中一个HASH分区的所有数据。通过`mr`函数在每个新的分区中计算95百分位的temperature,最后将结果合并汇总。 ``` // 创建数据库 deviceId = "device" + string(1..100000) db1 = database("", VALUE, 2019.06.01..2019.06.30) db2 = database("", HASH, SYMBOL:20) db = database("dfs://repartitionExample", COMPO, [db1, db2]) // 创建DFS表 t = db.createPartitionedTable(table(100000:0, `deviceId`time`temperature, [SYMBOL,DATETIME,DOUBLE]), `tb, `time`deviceId) n = 3000000 t.append!(table(rand(deviceId, n) as deviceId, 2019.06.01T00:00:00 + rand(86400 * 10, n) as time, 60 + norm(0.0, 5.0, n) as temperature)) // 重新分区 ds = repartitionDS() ds.transDS!(def (mutable t) { update t set x0 = nullFill(x0, avg(x0)), x1 = nullFill(x1, avg(x1)), x2 = nullFill(x2, avg(x2)), x3 = nullFill(x3, avg(x3)) return t }) randomForestRegressor(ds, `y, `x0`x1`x2`x3) ``` 另一个转换数据源的例子是2.2节提到的逻辑回归的脚本实现。在2.2节的实现中,map函数调用中包含了从数据源的表中取出对应列,转换成矩阵的操作,这意味着每一轮迭代都会发生这些操作。而实际上,每轮迭代都会使用同样的输入矩阵,这个转换步骤只需要调用一次。因此,可以用`transDS!`将数据源转换成一个包含x, xt和y矩阵的三元组: ``` def myLrTrans(t, yColName, xColNames, intercept) { if (intercept) x = matrix(t[xColNames], take(1.0, t.rows())) else x = matrix(t[xColNames]) xt = x.transpose() y = t[yColName] return [x, xt, y] } def myLrMap(input, lastFinal) { x, xt, y = input placeholder, placeholder, theta = lastFinal // 之后的计算和2.2节相同 } // myLrFinal和myLrTerm函数和2.2节相同 def myLr(mutable ds, yColName, xColNames, intercept, initTheta, tol) { ds.transDS!(myLrTrans{, yColName, xColNames, intercept}) logLik, grad, theta = imr(ds, [0, 0, initTheta], myLrMap, +, myLrFinal, myLrTerm{, , tol}) return theta } ```