PollingClient
PollingClient 在订阅数据时,会返回一个消息队列,用户可以从该消息队列中获取和处理数据。本小节将从构造函数、订阅、读取数据和取消订阅对 PollingClient 进行介绍,并展示一个使用示例。
构造方法
PollingClient(int listeningPort = 0)
其中,listeningPort 表示单线程客户端的订阅端口号。
PollingClient(StreamingClientConfig config);
其中: StreamingClientConfig 为流订阅 API 的配置,当前定义如下:
enum class TransportationProtocol {
TCP, UDP,
};
enum class SubscribeState {
Connected, // 连接成功,不代表有数据传来
Disconnected, // 连接断开
Resubscribing, // 正在重试
};
struct SubscribeInfo {
std::string hostName;
int port;
std::string tableName;
std::string actionName;
};
using SubscribeCallbackT = std::function<bool(const SubscribeState state, const SubscribeInfo &info)>;
struct StreamingClientConfig {
TransportationProtocol protocol{TransportationProtocol::TCP};
SubscribeCallbackT callback;
int netTimeout{30000};
};
protocol 为传输协议 :
-
为 TCP (默认值)时,与 3.00.2 及之前版本功能相同。
-
为 UDP 时,通过 Aeron 库与 DolphinDB Server 进行 UDP 组播通信。
使用 UDP 协议时:
-
当前 UDP 通信方式仅支持订阅接口 1,且不支持 resub, filter, backupSites, resubscribeInterval, subOnce 等参数。
-
C++ API 仅支持在 Linux 下编译和运行。
-
C++ API 在编译时依赖于 Aeron 库,但对 Aeron 库的版本没有限制。
-
C++ API 以嵌入式 Aeron Media Driver 的方式运行,客户无需单独启动 Aeron 线程。
-
当 C++ API 使用 UDP 进行流订阅时,会在
/dev/shm目录下创建一个以dolphindb_udp_<pid>命名的文件夹供 Aeron 使用,订阅结束后会自动删除。如客户程序异常退出,请手动删除相关文件夹。 -
Server 须为 Linux 3.00.0 以上版本。
-
Server 每次发布的信息大小上限为 2MB。如果订阅未收到数据,请检查是否超过了该限制。如有需要,可以调整配置项
maxMsgNumPerBlock以解决此问题。
callback 为状态发生变化时的回调函数,请参考以下示例:
int main(){
auto handler = [](const Message &msg) {
std::cout << "this is a message" << std::endl;
};
auto stateCallback = [](const SubscribeState state, const SubscribeInfo &info) {
std::cout << "state: " << static_cast<int>(state)
<< ", table: " << info.tableName << std::endl;
return true;
};
StreamingClientConfig config {
.callback = stateCallback,
};
PollingClient client(config);
client.subscribe("localhost", 8848, "shared_stream_table", "action1", -1, true, nullptr, false, false, "admin", "123456");
std::this_thread::sleep_for(100s);
client.unsubscribe("localhost", 8848, "shared_stream_table", "action1");
return 0;
}
netTimeout 控制流订阅的网络异常检测和建连超时,单位为毫秒,默认值为 30,000。若设置为 0,保活检测使用 30
秒的默认值,且不额外设置 TCP 建连超时。
订阅
MessageQueueSP subscribe(
string host,
int port,
string tableName,
string actionName,
int64_t offset = -1,
bool resub = true,
const VectorSP &filter = nullptr,
bool msgAsTable = false,
bool allowExists = false,
string userName="",
string password="",
const StreamDeserializerSP &blobDeserializer = nullptr,
const std::vector<std::string>& backupSites = std::vector<std::string>(),
int resubscribeInterval = 100,
bool subOnce = false,
int resubscribeTimeout = 0
)
参数
host:表示发布端节点的主机名。port:表示发布端节点的端口号。tableName:表示发布端上共享流数据表的名称。actionName:表示订阅任务的名称。传入值可以包含字母、数字和下划线。offset:表示订阅任务开始后的第一条消息所在的位置。消息是流数据表中的行。如果没有指定 offset、或其为负数、或超过了流数据表的记录行数,则订阅将会从流数据表的当前行开始。offset 与流数据表创建时的第一行对应。如果某些行因为内存限制被删除,在决定订阅开始的位置时,这些行仍然考虑在内。resub:表示订阅中断后,API 是否会自动重订阅。设置为 true 时,建议同时配置 resubscribeInterval 和 resubscribeTimeout,以防止长期频繁发起重订阅而占用过多系统资源。filter:一个向量,表示过滤条件。流数据表过滤列在 filter 中的数据才会发布到订阅端,不在 filter 中的数据不会发布。msgAsTable:只有设置了 batchSize 参数,该参数才会生效。若其设置为 true,则订阅的数据会转换为 Table;若其设置为 false,则订阅的数据会转换成 AnyVector。allowExists:若该参数设置为 true,则支持对同一个订阅流表使用多个 handler 处理数据;若其设置为 False,则不支持同一个订阅流表使用多个 handler 处理数据。userName:表示 API 所连接服务器的登录用户名。password:表示 API 所连接服务器的登录密码。blobDeserializer:表示订阅的异构流表对应的反序列化器。backupSites: 字符串列表,可选,表示备用的发布端节点列表,由节点 ip 和端口号组成,例如 [“192.168.0.1:8848“, “192.168.0.2:8849“]。- 指定backupSites,表示启动主备节点切换。如果发生节点切换(例如连接断开),会在可用节点列表中不断轮询订阅,直至订阅成功。
- 若配置该参数,用户需保证主节点(由 host 和 port 参数指定)和备用节点上存在相同结构、相同数据的同名流数据表。否则,可能出现订阅的数据不符合预期。
- 若订阅的是高可用流表,重连时将从 backupSites 指定的节点列表中选择节点。
- 取消订阅需使用主节点 ip 和端口号进行取消。
- resubscribeInterval:一个非负整数,表示在断连后每次尝试进行重新订阅的最小时间间隔(单位毫秒),默认值为 100。
subOnce:布尔值,False 表示在节点切换时,不去尝试之前成功连接过的节点。默认值为 False。- resubscribeTimeout:整型,表示重订阅的超时时间(单位毫秒)。默认值为0,表示无限重试。建议配置 resub=true 时,同时配置该参数,以防止重订阅过程耗时太长而占用过多系统资源。
返回
该函数返回指向消息队列的指针。
读取数据
MessageQueue 类实现了一个队列数据结构用于存储 Message 类实例。当一个
PollingClient 对象创建订阅时会返回一个 MessageQueue
对象。MessageQueue 支持以下方法从队列中读取数据:
poll
从队列中取出一条消息。如果当前队列中没有数据,则等待 milliSeconds。
bool poll(Message &item, int milliSeconds)- item 取出的消息。
- milliSeconds 最长等待时间,单位是毫秒。
- true 表示本次读取数据成功。
- false 表示在最长等待时间内队列中始终没有数据。
pop
从队列中取出一条消息。如果当前队列中没有数据,则将持续等待直到队列收到数据。
void pop(Message &item)- item 取出的消息。
取消订阅
bool unsubscribe(
string host,
int port,
string tableName,
string actionName
)
参数值需要与订阅时填入参数值相同。
取消订阅成功时返回 true,否则返回 false。
使用示例
DolphinDB 脚本:新建一个Stream表,然后共享该表,名为shared_stream_table。
rt = streamTable(`XOM`GS`AAPL as id, 102.1 33.4 73.6 as x)
share rt as shared_stream_table
C++代码:
本例使用2.00.10版本的DolphinDB server,在订阅一秒后取消订阅
int main(int argc, const char **argv)
{
PollingClient client(0);
MessageQueueSP queue = client.subscribe("127.0.0.1", 8848, "shared_stream_table", "action3", 0);
ThreadSP t = new Thread(new Executor([&client](){
Util::sleep(1000);
client.unsubscribe("127.0.0.1", 8848, "shared_stream_table", "action3");
}));
t->start();
Message msg;
while (queue->poll(msg, 1000) && !msg.isNull()) {
std::cout << msg->getString() << std::endl;
}
}
示例:将 TCP 流订阅的网络超时配置为 10 秒。
#include "Streaming.h"
using namespace dolphindb;
StreamingClientConfig config;
config.netTimeout = 10'000; // 单位:毫秒
PollingClient client(config);
