# DolphinDB 插件开发教程 DolphinDB 支持动态加载外部插件,以扩展系统功能。插件用 C\+\+ 编写,需要编译成 ".so" 或 ".dll" 共享库文件。插件使用的流程请参考 DolphinDB Plugin 的 [GitHub 页面](https://github.com/dolphindb/DolphinDBPlugin)。本文着重介绍如何开发插件,并详细介绍以下几个具体场景的插件开发流程: - [1. 如何开发插件](#1-如何开发插件) - [1.1 基本概念](#11-基本概念) - [1.2 创建变量](#12-创建变量) - [1.3 异常处理和参数校验](#13-异常处理和参数校验) - [1.3.1 异常处理](#131-异常处理) - [1.3.2 参数校验的范例](#132-参数校验的范例) - [1.4 调用 DolphinDB 内置函数](#14-调用-dolphindb-内置函数) - [2. 如何开发支持时间序列数据处理的插件函数](#2-如何开发支持时间序列数据处理的插件函数) - [3. 如何开发用于处理分布式 SQL 的聚合函数](#3-如何开发用于处理分布式-sql-的聚合函数) - [3.1 聚合函数范例](#31-聚合函数范例) - [3.2 在 DolphinDB 中调用函数](#32-在-dolphindb-中调用函数) - [3.3 随机访问大数组](#33-随机访问大数组) - [3.4 应该选择哪种方式访问向量](#34-应该选择哪种方式访问向量) - [4. 如何开发支持新的分布式算法的插件函数](#4-如何开发支持新的分布式算法的插件函数) - [4.1 分布式算法范例](#41-分布式算法范例) - [4.2 在 DolphinDB 中调用函数](#42-在-dolphindb-中调用函数) - [5. 如何开发支持流数据处理的插件函数](#5-如何开发支持流数据处理的插件函数) - [6. 如何开发支持外部数据源的插件函数](#6-如何开发支持外部数据源的插件函数) - [6.1 数据格式描述](#61-数据格式描述) - [6.2 `extractMyDataSchema` 函数](#62-extractmydataschema-函数) - [6.3 `loadMyData` 函数](#63-loadmydata-函数) - [6.4 `loadMyDataEx` 函数](#64-loadmydataex-函数) - [6.5 `myDataDS` 函数](#65-mydatads-函数) - [7. 如何在插件代码中构造并使用 sql 语句](#7-如何在插件代码中构造并使用-sql-语句) - [7.1 使用步骤](#71-使用步骤) - [7.1.1 引入 *ScalarImp.h* 头文件](#711-引入-scalarimph-头文件) - [7.1.2 将待查询的 Table 对象放入 Heap 中](#712-将待查询的-table-对象放入-heap-中) - [7.1.3 拼接 sql 字符串](#713-拼接-sql-字符串) - [7.1.4 执行 sql](#714-执行-sql) - [7.2 完整代码示例](#72-完整代码示例) - [7.2.1 `select * from t` 的完整代码示例](#721-select--from-t-的完整代码示例) - [7.2.2 `select avg(x) from t` 的完整代码示例](#722-select-avgx-from-t-的完整代码示例) - [8. 常见问题](#8-常见问题) - [8.1 如何处理开发的 windows 版本插件加载时的错误提示:"The specified module could not be found"?](#81-如何处理开发的-windows-版本插件加载时的错误提示the-specified-module-could-not-be-found) - [8.2 插件开发时需要包含哪些库和头文件?](#82-插件开发时需要包含哪些库和头文件) - [8.3 编译时需要包含哪些选项?](#83-编译时需要包含哪些选项) - [8.4 如何处理编译时出现包含 std::\_\_cxx11 字样的链接问题(undefined reference)?](#84-如何处理编译时出现包含-std__cxx11-字样的链接问题undefined-reference) - [8.5 如何加载插件,可以卸载后重新加载吗?](#85-如何加载插件可以卸载后重新加载吗) - [8.6 如何处理执行插件函数时的报错信息:"Connnection refused:connect" 或节点 crash 问题?](#86-如何处理执行插件函数时的报错信息connnection-refusedconnect-或节点-crash-问题) - [8.7 如何处理执行插件函数时的错误提示:"Cannot recognize the token xxx"?](#87-如何处理执行插件函数时的错误提示cannot-recognize-the-token-xxx) - [9. 附件](#9-附件) # 1. 如何开发插件 ## 1.1 基本概念 DolphinDB 的插件实现了能在脚本中调用的函数。一个插件函数可能是运算符函数(operator function),也可能是系统函数(system function),它们的区别在于,前者接受的参数个数小于等于 2,而后者可以接受任意个参数,并支持会话的访问操作。 开发一个运算符函数,需要编写一个原型为 ConstantSP (const ConstantSP& a, const ConstantSP& b) 的 C\+\+ 函数。当函数参数个数为 2 时,a 和 b 分别为插件函数的第一和第二个参数;当参数个数为 1 时,b 是一个占位符,没有实际用途;当没有参数时,a 和 b 均为占位符。 开发一个系统函数,需要编写一个原型为 ConstantSP (Heap* heap, vector& args) 的 C\+\+ 函数。用户在 DolphinDB 中调用插件函数时传入的参数,都按顺序保存在 C++ 的向量 args 中。heap 参数不需要用户传入。 函数原型中的 ConstantSP 可以表示绝大多数 DolphinDB 对象(标量、向量、矩阵、表,等等)。其他常用的派生自它的变量类型有 VectorSP(向量)以及 TableSP(表)等。 开发插件时所需的头文件通常位于 DolphinDBPlugin 项目的 include 目录下,用户可切换至指定版本下载。 此处给出 Gitee 中某版本插件头文件的[链接](https://gitee.com/dolphindb/DolphinDBPlugin/tree/release200.10/include)。更多说明可参阅本文 8.常见问题。 ## 1.2 创建变量 创建标量,可以直接用 new 语句创建头文件 ScalarImp.h 中声明的类型对象,并将它赋值给一个 ConstantSP。 ConstantSP 是一个经过封装的智能指针,会在变量的引用计数为 0 时自动释放内存,因此,用户不需要手动删除已经创建的变量: ```cpp ConstantSP i = new Int(1); // 相当于 1i ConstantSP d = new Date(2019, 3, 14); // 相当于 2019.03.14 ConstantSP s = new String("DolphinDB"); // 相当于 "DolphinDB" ConstantSP voidConstant = new Void(); // 创建一个 void 类型变量,常用于表示空的函数参数 ``` 头文件 Util.h 声明了一系列函数,用于快速创建某个类型和格式的变量: ```cpp VectorSP v = Util::createVector(DT_INT, 10); // 创建一个初始长度为 10 的 INT 类型向量 v->setInt(0, 60); // 相当于 v[0] = 60 VectorSP t = Util::createVector(DT_ANY, 0); // 创建一个初始长度为 0 的 ANY 类型向量(元组) t->append(new Int(3)); // 相当于 t.append!(3) t->get(0)->setInt(4); // 相当于 t[0] = 4 // 这里不能用 t->setInt(0, 4),因为 t 是一个元组,setInt(0, 4) 只对 INT 类型的向量有效 ConstantSP seq = Util::createIndexVector(5, 10); // 相当于 5..14 int seq0 = seq->getInt(0); // 相当于 seq[0] ConstantSP mat = Util::createDoubleMatrix(5, 10);// 创建一个 10 行 5 列的 DOUBLE 类型矩阵 mat->setColumn(3, seq); // 相当于 mat[3] = seq ``` ## 1.3 异常处理和参数校验 ### 1.3.1 异常处理 插件开发时的异常抛出和处理,和一般 C\+\+ 开发中一样,都通过 throw 关键字抛出异常, try 语句块处理异常。DolphinDB 在头文件 Exceptions.h 中声明了异常类型。插件函数若遇到运行时错误,一般抛出 RuntimeException。 在插件开发时,通常会校验函数参数,如果参数不符合要求,抛出一个 IllegalArgumentException。常用的参数校验函数有: - `ConstantSP->getType()`:返回变量的类型(INT, CHAR, DATE 等等),DolphinDB 的类型定义在头文件 Types.h 中。 - `ConstantSP->getCategory()`:返回变量的类别,常用的类别有 INTEGRAL(整数类型,包括 INT, CHAR, SHORT, LONG 等)、FLOATING(浮点数类型,包括 FLOAT, DOUBLE 等)、TEMPORAL(时间类型,包括 TIME, DATE, DATETIME 等)、LITERAL(字符串类型,包括 STRING, SYMBOL 等),都定义在头文件 Types.h 中。 - `ConstantSP->getForm()`:返回变量的格式(标量、向量、表等等),DolphinDB 的格式定义在头文件 Types.h 中。 - `ConstantSP->isVector()`:判断变量是否为向量。 - `ConstantSP->isScalar()`:判断变量是否为标量。 - `ConstantSP->isTable()`:判断变量是否为表。 - `ConstantSP->isNumber()`:判断变量是否为数字类型。 - `ConstantSP->isNull()`:判断变量是否为空值。 - `ConstantSP->getInt()`:获得变量对应的整数值。 - `ConstantSP->getString()`:获得变量对应的字符串。 - `ConstantSP->size()`:获得变量的长度。 更多参数校验函数一般在头文件 CoreConcept.h 的 Constant 类方法中。 ### 1.3.2 参数校验的范例 本节将开发一个插件函数用于求非负整数的阶乘,返回一个 LONG 类型变量。DolphinDB 中 LONG 类型的最大值为 2^63-1,能表示的阶乘最大为 25!,因此只有 0~25 范围内的参数是合法的。 ```cpp #include "CoreConcept.h" #include "Exceptions.h" #include "ScalarImp.h" ConstantSP factorial(const ConstantSP &n, const ConstantSP &placeholder) { string syntax = "Usage: factorial(n)."; if (!n->isScalar() || n->getCategory() != INTEGRAL) throw IllegalArgumentException("factorial", syntax + "n must be an integral scalar."); int nValue = n->getInt(); if (nValue < 0 || nValue> 25) throw IllegalArgumentException("factorial", syntax + "n must be a non-negative integer less than 26."); long long fact = 1; for (int i = nValue; i> 0; i--) fact *= i; return new Long(fact); } ``` ## 1.4 调用 DolphinDB 内置函数 有时需要调用 DolphinDB 的内置函数对数据进行处理。有些类已经定义了部分常用的内置函数作为方法: ```cpp VectorSP v = Util::createIndexVector(1, 100); ConstantSP avg = v->avg(); // 相当于 avg(v) ConstantSP sum2 = v->sum2(); // 相当于 sum2(v) v->sort(false); // 相当于 sort(v, false) ``` 如果需要调用其它内置函数,插件函数的类型必须是系统函数。通过 heap->currentSession()->getFunctionDef 函数获得一个内置函数,然后用 `call` 方法调用它。如果该内置函数是运算符函数,应调用原型 call(Heap, const ConstantSP&, const ConstantSP&);如果是系统函数,应调用原型 call(Heap, vector&)。以下是调用内置函数 `cumsum` 的一个例子: ```cpp ConstantSP v = Util::createIndexVector(1, 100); v->setTemporary(false); //v 的值可能在内置函数调用时被修改。如果不希望它被修改,应先调用 setTemporary(false) FunctionDefSP cumsum = heap->currentSession()->getFunctionDef("cumsum"); ConstantSP result = cumsum->call(heap, v, new Void()); // 相当于 cumsum(v),这里的 new Void() 是一个占位符,没有实际用途 ``` # 2. 如何开发支持时间序列数据处理的插件函数 DolphinDB 的特色之一在于它对时间序列有良好支持。本章以编写一个 [msum](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/m/msum.html) 函数的插件为例,介绍如何开发插件函数支持时间序列数据处理。 时间序列处理函数通常接受向量作为参数,并对向量中的每个元素进行计算处理。在本例中,`msum` 函数接受两个参数:一个向量和一个窗口大小。它的原型是: ```cpp ConstantSP msum(const ConstantSP &X, const ConstantSP &window); ``` `msum` 函数的返回值是一个和输入向量同样长度的向量。本例为简便起见,假定返回值是一个 DOUBLE 类型的向量。可以通过 Util::createVector 函数预先为返回值分配空间: ```cpp int size = X->size(); int windowSize = window->getInt(); ConstantSP result = Util::createVector(DT_DOUBLE, size); ``` 在 DolphinDB 插件编写时处理向量,可以循环使用 `getDoubleConst`, `getIntConst` 等函数,批量获得一定长度的只读数据,保存在相应类型的缓冲区中,从缓冲区中取得数据进行计算。这样做的效率比循环使用 `getDouble`, `getInt` 等函数要高。本例为简便起见,统一使用 `getDoubleConst`,每次获得长度为 Util::BUF_SIZE 的数据。这个函数返回一个 const double* ,指向缓冲区头部: ```cpp double buf[Util::BUF_SIZE]; INDEX start = 0; while (start < size) { int len = std::min(Util::BUF_SIZE, size - start); const double *p = X->getDoubleConst(start, len, buf); for (int i = 0; i < len; i++) { double val = p[i]; // ... } start += len; } ``` 在本例中,`msum` 将计算 X 中长度为 windowSize 的窗口中所有数据的和。可以用一个临时变量 tmpSum 记录当前窗口的和,每当窗口移动时,只要给 tmpSum 增加新窗口尾部的值,减去旧窗口头部的值,就能计算得到当前窗口中数据的和。为了将计算值写入 result,可以循环用 result->getDoubleBuffer 获取一个可读写的缓冲区,写完后使用 result->setDouble 函数将缓冲区写回数组。`setDouble` 函数会检查给定的缓冲区地址和变量底层储存的地址是否一致,如果一致就不会发生数据拷贝。在多数情况下,用 `getDoubleBuffer` 获得的缓冲区就是变量实际的存储区域,这样能减少数据拷贝,提高性能。 需要注意的是,DolphinDB 用 DOUBLE 类型的最小值(已经定义为宏 DBL_NMIN )表示 DOUBLE 类型的 NULL 值,要专门判断。 返回值的前 windowSize - 1 个元素为 NULL。可以对 X 中的前 windowSize 个元素和之后的元素用两个循环分别处理,前一个循环只计算累加,后一个循环执行加和减的操作。最终的实现如下: ```cpp ConstantSP msum(const ConstantSP &X, const ConstantSP &window) { INDEX size = X->size(); int windowSize = window->getInt(); ConstantSP result = Util::createVector(DT_DOUBLE, size); double buf[Util::BUF_SIZE]; double windowHeadBuf[Util::BUF_SIZE]; double resultBuf[Util::BUF_SIZE]; double tmpSum = 0.0; INDEX start = 0; while (start < windowSize) { int len = std::min(Util::BUF_SIZE, windowSize - start); const double *p = X->getDoubleConst(start, len, buf); double *r = result->getDoubleBuffer(start, len, resultBuf); for (int i = 0; i < len; i++) { if (p[i] != DBL_NMIN) // p[i] is not NULL tmpSum += p[i]; r[i] = DBL_NMIN; } result->setDouble(start, len, r); start += len; } result->setDouble(windowSize - 1, tmpSum); // 上一个循环多设置了一个 NULL,填充为 tmpSum while (start < size) { int len = std::min(Util::BUF_SIZE, size - start); const double *p = X->getDoubleConst(start, len, buf); const double *q = X->getDoubleConst(start - windowSize, len, windowHeadBuf); double *r = result->getDoubleBuffer(start, len, resultBuf); for (int i = 0; i < len; i++) { if (p[i] != DBL_NMIN) tmpSum += p[i]; if (q[i] != DBL_NMIN) tmpSum -= q[i]; r[i] = tmpSum; } result->setDouble(start, len, r); start += len; } return result; } ``` # 3. 如何开发用于处理分布式 SQL 的聚合函数 在 DolphinDB 中,SQL 的聚合函数通常接受一个或多个向量作为参数,最终返回一个标量。在开发聚合函数的插件时,需要了解如何访问向量中的元素。 DolphinDB 中的向量有两种存储方式。一种是常规数组,数据在内存中连续存储,另一种是 [大数组](https://www.dolphindb.cn/cn/help/DataTypesandStructures/DataForms/Vector/BigArray.html),其中的数据分块存储。 本章将以编写一个求 [几何平均数](https://en.wikipedia.org/wiki/Geometric_mean) 的函数为例,介绍如何开发聚合函数,重点关注数组中元素的访问。 ## 3.1 聚合函数范例 几何平均数 `geometricMean` 函数接受一个向量作为参数。为了防止溢出,一般采用其对数形式计算,即 ```cpp geometricMean([x1, x2, ..., xn]) = exp((log(x1) + log(x2) + log(x3) + ... + log(xn))/n) ``` 为了实现这个函数的分布式版本,可以先开发聚合函数插件 `logSum`,用以计算某个分区上的数据的对数和,然后用 defg 关键字定义一个 reduce 函数,用 mapr 关键字定义一个 MapReduce 函数。 在 DolphinDB 插件开发中,对数组的操作通常要考虑它是常规数组还是大数组。可以用 `isFastMode` 函数判断: ```cpp ConstantSP logSum(const ConstantSP &x, const ConstantSP &placeholder) { if (((VectorSP) x)->isFastMode()) { // ... } else { // ... } } ``` 如果是常规数组,它在内存中连续存储。可以用 `getDataArray` 函数获得它数据的指针。假定数据是以 DOUBLE 类型存储的: ```cpp if (((VectorSP) x)->isFastMode()) { int size = x->size(); double *data = (double *) x->getDataArray(); double logSum = 0; for (int i = 0; i < size; i++) { if (data[i] != DBL_NMIN) // is not NULL logSum += std::log(data[i]); } return new Double(logSum); } ``` 如果是大数组,它在内存中分块存储。可以用 `getSegmentSize` 获得每个块的大小;用 `getDataSegment` 获得首个块的地址,以返回一个二级指针,指向一个指针数组,这个数组中的每个元素指向每个块的数据数组: ```cpp // ... else { int size = x->size(); int segmentSize = x->getSegmentSize(); double **segments = (double **) x->getDataSegment(); INDEX start = 0; int segmentId = 0; double logSum = 0; while (start < size) { double *block = segments[segmentId]; int blockSize = std::min(segmentSize, size - start); for (int i = 0; i < blockSize; i++) { if (block[i] != DBL_NMIN) // is not NULL logSum += std::log(block[i]); } start += blockSize; segmentId++; } return new Double(logSum); } ``` 以上代码针对 DOUBLE 数据类型。在实际开发中,数组的底层存储不一定是 DOUBLE 类型,也可能涉及多种数据类型,此时可采用泛型编程。附件中有一个泛型编程的代码例子。 ## 3.2 在 DolphinDB 中调用函数 通常需要实现一个聚合函数的非分布式版本和分布式版本,系统会基于哪个版本更高效来选择调用这个版本。 在 DolphinDB 中定义非分布式的 `geometricMean` 函数: ```cpp def geometricMean(x) { return exp(logSum::logSum(x) \ count(x)) } ``` 然后通过定义 Map 和 Reduce 函数,最终用 `mapr` 定义分布式的版本: ```cpp def geometricMeanMap(x) { return logSum::logSum(x) } defg geometricMeanReduce(myLogSum, myCount) { return exp(sum(myLogSum) \ sum(myCount)) } mapr geometricMean(x) { geometricMeanMap(x), count(x) -> geometricMeanReduce } ``` 如果是在单机环境中执行这个函数,只需要在执行的节点上加载插件。如果有数据位于远程节点,需要在每一个远程节点加载插件。可以手动在每个节点执行 `loadPlugin` 函数,也可以用以下脚本快速在每个节点上加载插件: ```cpp each(rpc{, loadPlugin, pathToPlugin}, getDataNodes()) ``` 通过以下脚本创建一个分区表,验证函数: ```cpp db = database("", VALUE, 1 2 3 4) t = table(take(1..4, 100) as id, rand(1.0, 100) as val) t0 = db.createPartitionedTable(t, `tb, `id) t0.append!(t) select geometricMean(val) from t0 group by id; ``` ## 3.3 随机访问大数组 可以对大数组进行随机访问,但要经过下标计算。用 `getSegmentSizeInBit` 函数获得块大小的二进制位数,通过位运算获得块的偏移量和块内偏移量: ```cpp int segmentSizeInBit = x->getSegmentSizeInBit(); int segmentMask = (1 << segmentSizeInBit) - 1; double **segments = (double **) x->getDataSegment(); int index = 3000000; // 想要访问的下标 double result = segments[index>> segmentSizeInBit][index & segmentMask]; // ^ 块的偏移量 ^ 块内偏移量 ``` ## 3.4 应该选择哪种方式访问向量 上一章 [如何开发支持时间序列数据处理的插件函数](#2 - 如何开发支持时间序列数据处理的插件函数) 介绍了通过 `getDoubleConst`, `getIntConst` 等一组方法获得只读缓冲区,以及通过 `getDoubleBuffer`, `getIntBuffer` 等一组方法获得可读写缓冲区。这两种访问向量的方法在实际开发比较通用。 本章介绍了通过 `getDataArray` 和 `getDataSegment` 方法直接访问向量的底层存储。在某些特别的场合,例如明确知道数据存储在大数组中,且知道数据的类型,这种方法比较适合。 # 4. 如何开发支持新的分布式算法的插件函数 在 DolphinDB database 中,MapReduce 是执行分布式算法的通用计算框架。DolphinDB 提供了 [mr](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/m/mr.html) 函数和 [imr](https://www.dolphindb.cn/cn/help/FunctionsandCommands/FunctionReferences/i/imr.html) 函数,使用户能通过脚本实现分布式算法。在编写分布式算法的插件时,使用的同样是这两个函数。对通用计算的详细介绍,可以参考[通用计算教程](general_computing.md)。本章主要介绍如何用 C\+\+ 语言编写自定义的 map, reduce 等函数,并调用 `mr` 和 `imr` 这两个函数,最终实现分布式计算。 ## 4.1 分布式算法范例 本章将使用 `mr`,实现一个函数,求分布式表中多个指定列中所有数据的平均值。我们会介绍编写 DolphinDB 分布式算法插件的整体流程,及需要注意的技术细节。 在插件开发中,用户自定义的 map, reduce, final, term 函数,可以是运算符函数,也可以是系统函数。 本例的 map 函数,对表的每个分区内的所有指定列做计算。每个分区返回一个长度为 2 的元组,包含数据之和,以及非空元素的个数。具体实现如下: ```cpp ConstantSP columnAvgMap(Heap *heap, vector &args) { TableSP table = args[0]; ConstantSP colNames = args[1]; double sum = 0.0; int count = 0; for (int i = 0; i < colNames->size(); i++) { string colName = colNames->getString(i); VectorSP col = table->getColumn(colName); sum += col->sum()->getDouble(); count += col->count(); } ConstantSP result = Util::createVector(DT_ANY, 2); result->set(0, new Double(sum)); result->set(1, new Int(count)); return result; } ``` 本例的 reduce 函数,是对 map 结果的相加。可使用 DolphinDB 的内置函数 `add`。用 heap->currentSession()->getFunctionDef("add") 获得这个函数: ```cpp FunctionDefSP reduceFunc = heap->currentSession()->getFunctionDef("add"); ``` 本例的 final 函数,是对 reduce 结果中的数据总和 `sum` 和非空元素个数 `count` 做除法,求得所有分区中对应列的平均数。具体实现如下: ```cpp ConstantSP columnAvgFinal(const ConstantSP &result, const ConstantSP &placeholder) { double sum = result->get(0)->getDouble(); int count = result->get(1)->getInt(); return new Double(sum / count); } ``` 定义了 map, reduce, final 等函数后,将它们导出为插件函数(在头文件的函数声明前加上 extern "C" ,并在加载插件的文本文件中列出这些函数),然后通过 heap->currentSession->getFunctionDef 获取这些函数,就能以这些函数为参数调用 `mr` 函数。如: ```cpp FunctionDefSP mapFunc = Heap->currentSession()->getFunctionDef("columnAvg::columnAvgMap"); ``` 在本例中,map 函数接受两个参数 table 和 colNames ,但 `mr` 只允许 map 函数有一个参数,因此需要以 [部分应用](https://www.dolphindb.cn/cn/help/Functionalprogramming/PartialApplication.html) 的形式调用 map 函数,可以用 Util::createPartialFunction 将它包装为部分应用,实现如下: ```cpp vector mapWithColNamesArgs {new Void(), colNames}; FunctionDefSP mapWithColNames = Util::createPartialFunction(mapFunc, mapWithColNamesArgs); ``` 用 heap->currentSession()->getFunctionDef("mr") 获得系统内置函数 `mr`,调用 mr->call 方法,就相当于在 DolphinDB 脚本中调用 `mr` 函数。 最后实现的 columnAvg 函数定义如下: ```cpp ConstantSP columnAvg(Heap *heap, vector &args) { ConstantSP ds = args[0]; ConstantSP colNames = args[1]; FunctionDefSP mapFunc = heap->currentSession()->getFunctionDef("columnAvg::columnAvgMap"); vector mapWithColNamesArgs = {new Void(), colNames}; FunctionDefSP mapWithColNames = Util::createPartialFunction(mapFunc, mapWithColNamesArgs); // columnAvgMap{, colNames} FunctionDefSP reduceFunc = heap->currentSession()->getFunctionDef("add"); FunctionDefSP finalFunc = heap->currentSession()->getFunctionDef("columnAvg::columnAvgFinal"); FunctionDefSP mr = heap->currentSession()->getFunctionDef("mr"); // mr(ds, columnAvgMap{, colNames}, add, columnAvgFinal) vector mrArgs = {ds, mapWithColNames, reduceFunc, finalFunc}; return mr->call(heap, mrArgs); } ``` ## 4.2 在 DolphinDB 中调用函数 如果是在单机环境中执行这个函数,只需要在执行的节点上加载插件。但如果计算需要用到位于远程节点的数据,就需要在每一个远程节点加载插件。可以手动在每个节点执行 `loadPlugin` 函数,也可以用以下脚本快速在每个节点上加载插件: ```cpp each(rpc{, loadPlugin, pathToPlugin}, getDataNodes()) ``` 加载插件后,用 `sqlDS` 函数生成数据源,并调用函数: ```cpp n = 100 db = database("dfs://testColumnAvg", VALUE, 1..4) t = db.createPartitionedTable(table(10:0, `id`v1`v2, [INT,DOUBLE,DOUBLE]), `t, `id) t.append!(table(take(1..4, n) as id, rand(10.0, n) as v1, rand(100.0, n) as v2)) ds = sqlDS(