Appearance
C++ Arrow 查询示例
下面示例是可独立复制的 C++17 客户端。它使用 libcurl 请求 /api/v1/query,再用 Apache Arrow C++ 的 arrow::ipc::RecordBatchStreamReader 读取 Arrow IPC Stream。
依赖与构建
需要 CMake 3.16+、libcurl、Apache Arrow C++ 开发包;检测到 OpenMP 时会启用并行归约,否则自动使用串行循环。下面示例计算交易中常见的成交量加权平均价(VWAP),Windows 可用 vcpkg:
powershell
vcpkg install curl:x64-windows arrow:x64-windows
cmake -S . -B build -DCMAKE_TOOLCHAIN_FILE=C:/src/vcpkg/scripts/buildsystems/vcpkg.cmake
cmake --build build --config ReleaseLinux/macOS:
bash
cmake -S . -B build -DCMAKE_BUILD_TYPE=Release
cmake --build build --parallel运行与环境变量
无需在本地启动服务,直接连接 FinORM 官方服务。请先准备有权限访问数据的授权 token,并确认目标表包含可计算的价格、数量列。示例只从环境变量读取配置,不会打印 token:
bash
export FINORM_ENDPOINT=http://finorm.zltquant.com:7761
export FINORM_TABLE=cneqa_mkt_dly
export FINORM_LIMIT=1000
export FINORM_PRICE_FIELD=price
export FINORM_VOLUME_FIELD=volume
export FINORM_TOKEN=your-token # 必填:替换为你的授权 token
./build/finorm_arrow_demoWindows PowerShell 使用 $env:FINORM_ENDPOINT="http://finorm.zltquant.com:7761" 等方式设置。程序会自动补全 /api/v1/query 路径;FINORM_PRICE_FIELD 与 FINORM_VOLUME_FIELD 可替换为实际字段名。FINORM_TOKEN 必须设置,不要把 token 写进源码、命令仓库或日志。
完整 C++ 示例
将下面内容保存为 main.cpp。示例只依赖公开 HTTP 契约,不读取本地数据文件,也不依赖服务端实现细节。
示例计算的是交易分析中常见的成交量加权平均价(VWAP):
text
VWAP = Σ(成交价 × 成交量) / Σ(成交量)与传统的逐行 JSON 处理相比,Arrow 直接提供列式数组,程序可以一次读取价格列和成交量列,在每个 RecordBatch 内并行归约。传统方式适合小结果和人工查看,优点是结构直观;Arrow 方式适合大结果和数值计算,优点是减少 JSON 解析与对象创建开销,便于 OpenMP、SIMD 或后续 GPU 计算。两种方式使用相同的查询条件和权限,计算结果应保持一致。
cpp
// 读取 Arrow IPC Stream,计算成交量加权平均价(VWAP)。
// 代码只依赖公开 HTTP 接口,不依赖服务端内部类或本地数据文件。
#include <arrow/io/api.h>
#include <arrow/ipc/api.h>
#include <arrow/type.h>
#include <curl/curl.h>
#include <cstdlib>
#include <iostream>
#include <memory>
#include <sstream>
#include <string>
#ifdef _OPENMP
#include <omp.h>
#endif
namespace {
struct HttpResponse {
long status = 0;
std::string body;
std::string content_type;
};
// 从环境变量读取配置,避免把服务地址和授权 token 写进源码。
std::string env(const char* key, const char* fallback) {
const char* value = std::getenv(key);
return value && *value ? value : fallback;
}
// libcurl 回调:把 HTTP 响应体原样保存为 Arrow IPC 字节流。
size_t write_body(char* data, size_t size, size_t count, void* target) {
static_cast<std::string*>(target)->append(data, size * count);
return size * count;
}
// libcurl 回调:读取响应媒体类型,确认服务端返回 Arrow Stream。
size_t write_header(char* data, size_t size, size_t count, void* target) {
std::string line(data, size * count);
const auto colon = line.find(':');
if (colon != std::string::npos && line.substr(0, colon) == "Content-Type") {
auto* response = static_cast<HttpResponse*>(target);
response->content_type = line.substr(colon + 1);
}
return size * count;
}
// 表名来自环境变量,先做最小 JSON 转义,避免破坏请求体结构。
std::string json_escape(const std::string& value) {
std::string escaped;
for (char ch : value) {
if (ch == '\\' || ch == '"') escaped.push_back('\\');
escaped.push_back(ch);
}
return escaped;
}
// 发起一次 Arrow 查询:请求体声明 format=arrow,token 通过 Bearer 头传递。
bool query(const std::string& endpoint, const std::string& table,
const std::string& token, HttpResponse* response) {
CURL* curl = curl_easy_init();
if (!curl) return false;
// limit=1000 只是示例值,实际项目应结合 query_spec 和业务分页策略调整。
const std::string body = "{\"table\":\"" + json_escape(table) +
"\",\"format\":\"arrow\",\"limit\":1000}";
curl_slist* headers = nullptr;
headers = curl_slist_append(headers, "Content-Type: application/json");
headers = curl_slist_append(headers, "Accept: application/vnd.apache.arrow.stream");
std::string authorization;
authorization = "Authorization: Bearer " + token;
headers = curl_slist_append(headers, authorization.c_str());
// 同时声明 JSON 请求和 Arrow Stream 响应,错误响应仍由服务端返回标准 JSON。
curl_easy_setopt(curl, CURLOPT_URL, endpoint.c_str());
curl_easy_setopt(curl, CURLOPT_POST, 1L);
curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body.c_str());
curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, write_body);
curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response->body);
curl_easy_setopt(curl, CURLOPT_HEADERFUNCTION, write_header);
curl_easy_setopt(curl, CURLOPT_HEADERDATA, response);
const CURLcode code = curl_easy_perform(curl);
if (code == CURLE_OK) curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &response->status);
curl_slist_free_all(headers);
curl_easy_cleanup(curl);
return code == CURLE_OK;
}
// 把 Arrow 数值数组中的一个元素读成 double,便于演示 VWAP 公式。
// 生产计算如需无损 decimal,应保留 decimal 类型,不要直接转 double。
double number_at(const std::shared_ptr<arrow::Array>& array, int64_t index) {
switch (array->type_id()) {
case arrow::Type::DOUBLE:
return std::static_pointer_cast<arrow::DoubleArray>(array)->Value(index);
case arrow::Type::FLOAT:
return std::static_pointer_cast<arrow::FloatArray>(array)->Value(index);
case arrow::Type::INT64:
return static_cast<double>(std::static_pointer_cast<arrow::Int64Array>(array)->Value(index));
case arrow::Type::INT32:
return static_cast<double>(std::static_pointer_cast<arrow::Int32Array>(array)->Value(index));
default:
return 0.0;
}
}
// 在一个 RecordBatch 内计算两个可加的中间量:
// price_volume = Σ(price × volume)
// volume_total = Σ(volume)
// OpenMP 只写线程私有归约变量,避免多个线程竞争同一个累加器。
struct BatchSums {
double price_volume = 0.0;
double volume = 0.0;
};
BatchSums batch_sums(const std::shared_ptr<arrow::Array>& price,
const std::shared_ptr<arrow::Array>& volume) {
double price_volume = 0.0;
double volume_total = 0.0;
#ifdef _OPENMP
#pragma omp parallel for reduction(+ : price_volume, volume_total)
#endif
// 每一行对应一笔成交,空值行不参与 VWAP 分子和分母。
for (int64_t i = 0; i < price->length(); ++i) {
if (!price->IsNull(i) && !volume->IsNull(i)) {
const double price_value = number_at(price, i);
const double volume_value = number_at(volume, i);
price_volume += price_value * volume_value;
volume_total += volume_value;
}
}
return {price_volume, volume_total};
}
} // namespace
int main() {
// 官方服务地址是默认值,用户只需提供 token 和自己的表/字段配置。
std::string endpoint = env("FINORM_ENDPOINT", "http://finorm.zltquant.com:7761");
if (endpoint.find("/api/v1/query") == std::string::npos) {
while (!endpoint.empty() && endpoint.back() == '/') endpoint.pop_back();
endpoint += "/api/v1/query";
}
const std::string table = env("FINORM_TABLE", "your_table");
const std::string price_name = env("FINORM_PRICE_FIELD", "price");
const std::string volume_name = env("FINORM_VOLUME_FIELD", "volume");
const std::string token = env("FINORM_TOKEN", "");
if (token.empty()) {
std::cerr << "FINORM_TOKEN must be set\n";
return 1;
}
// 初始化 libcurl 全局资源,整个进程只执行一次。
if (curl_global_init(CURL_GLOBAL_DEFAULT) != CURLE_OK) return 1;
HttpResponse response;
// 下载完整 Arrow Stream 后再交给 Arrow IPC reader 解码。
const bool ok = query(endpoint, table, token, &response);
curl_global_cleanup();
if (!ok || response.status < 200 || response.status >= 300) {
std::cerr << "HTTP request failed: " << response.status << "\n" << response.body << "\n";
return 1;
}
// BufferReader 按内存中的顺序字节读取 IPC Stream,不需要 ARROW1 文件尾索引。
auto buffer = arrow::Buffer::FromString(response.body);
auto input = arrow::io::BufferReader::Make(buffer);
if (!input.ok()) { std::cerr << input.status().ToString() << "\n"; return 1; }
auto reader = arrow::ipc::RecordBatchStreamReader::Open(*input);
if (!reader.ok()) { std::cerr << reader.status().ToString() << "\n"; return 1; }
// 跨 RecordBatch 合并两个中间量,最后再套用 VWAP 公式。
double price_volume_total = 0.0;
double volume_total = 0.0;
int64_t rows = 0;
while (true) {
auto batch = (*reader)->ReadNext();
if (!batch.ok()) { std::cerr << batch.status().ToString() << "\n"; return 1; }
if (!*batch) break;
const int price_index = (*batch)->schema()->GetFieldIndex(price_name);
const int volume_index = (*batch)->schema()->GetFieldIndex(volume_name);
if (price_index < 0 || volume_index < 0) {
std::cerr << "missing numeric field in Arrow schema\n";
return 1;
}
// 只在当前 batch 生命周期内使用列数组,避免保存失效指针。
const BatchSums sums = batch_sums((*batch)->column(price_index), (*batch)->column(volume_index));
price_volume_total += sums.price_volume;
volume_total += sums.volume;
rows += (*batch)->num_rows();
}
if (volume_total == 0.0) {
std::cerr << "cannot calculate VWAP: total volume is zero\n";
return 1;
}
// VWAP = Σ(price × volume) / Σ(volume)。
const double vwap = price_volume_total / volume_total;
std::cout << "rows=" << rows << " vwap=" << vwap
<< " (sum_price_volume=" << price_volume_total
<< ", sum_volume=" << volume_total << ")\n";
return 0;
}将下面内容保存为同目录的 CMakeLists.txt:
cmake
cmake_minimum_required(VERSION 3.16)
project(finorm_arrow_demo LANGUAGES CXX)
set(CMAKE_CXX_STANDARD 17)
set(CMAKE_CXX_STANDARD_REQUIRED ON)
find_package(CURL REQUIRED)
find_package(Arrow CONFIG REQUIRED)
find_package(OpenMP QUIET)
add_executable(finorm_arrow_demo main.cpp)
if(TARGET Arrow::arrow_shared)
target_link_libraries(finorm_arrow_demo PRIVATE CURL::libcurl Arrow::arrow_shared)
elseif(TARGET Arrow::arrow)
target_link_libraries(finorm_arrow_demo PRIVATE CURL::libcurl Arrow::arrow)
else()
target_link_libraries(finorm_arrow_demo PRIVATE CURL::libcurl Arrow::arrow_static)
endif()
if(OpenMP_CXX_FOUND)
target_link_libraries(finorm_arrow_demo PRIVATE OpenMP::OpenMP_CXX)
endif()Linux/macOS 示例编译命令:
bash
c++ -std=c++17 main.cpp -o finorm_arrow_demo \
$(pkg-config --cflags --libs arrow) -lcurl -fopenmp
./finorm_arrow_demo如果编译器没有 OpenMP,去掉 -fopenmp 即可,程序会自动使用串行循环。
请求、读取与计算
客户端设置 Content-Type: application/json 和 Accept: application/vnd.apache.arrow.stream,POST body 中包含表名、字段、行数上限及 format: "arrow"。响应必须是 2xx;非 2xx 直接输出标准错误体并以非零状态退出。响应头可用于记录 X-FinORM-Returned-Rows、X-FinORM-Truncated、X-FinORM-Request-Id 等审计信息。
Arrow 响应是 IPC Stream(不是带 ARROW1 文件头的 IPC File),程序按顺序读取每个 RecordBatch,计算常见的成交量加权平均价(VWAP)。对每一行,先累加 price × volume 和 volume,所有批次合并后计算:VWAP = Σ(price × volume) / Σ(volume)。这对应交易分析中“某段时间内实际成交的平均价格”,比简单算术平均更能反映大成交量订单的影响。整数、浮点和 decimal128 均可参与示例计算,最终示例汇总为 double;若业务需要保留高精度,应在进入计算内核前使用 Arrow decimal 类型或显式转换策略。
OpenMP 路径使用两个循环归约变量分别累加 Σ(price × volume) 与 Σ(volume),线程不竞争共享总和,batch 完成后再汇总;未找到 OpenMP 时仍可编译运行,只是使用串行循环。相同模式可扩展到收益率、成交额、波动率等按行可加的指标:先计算线程局部中间量,最后做一次合并。Arrow 数组只在当前 RecordBatch 生命周期内有效,不要在 batch 释放后保存裸指针。
相关格式与安全提示
跨语言消费或需要保留高精度 DECIMAL 时优先 Arrow;纯 Python 且直接喂 NumPy/numba 时可选 npz。查询仍须遵循公开的 Catalog、query_spec、最大行数和时间范围约束,不接受任意 SQL。调用官方服务时必须携带有效的 Bearer token。