ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

StarRocks Stream Load 实战指南:从一条 curl 到生产级实时数据导入

StarRocks Stream Load 实战指南:从一条 curl 到生产级实时数据导入 StarRocks Stream Load 实战指南从一条 curl 到生产级实时数据导入【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocksStarRocks 的 Stream Load 让数据导入回归到一次 HTTP 请求的粒度文件 PUT 上去落表即可查询无需额外的调度层。实测单节点可维持每分钟数十万行的导入吞吐这也是它成为传感器遥测、日志分析等场景默认实时链路的原因。本文按先跑通 → 再跑对 → 然后跑快 → 最后跑稳的顺序展开涉及的都是可以直接复制执行的命令和参数。一次请求如何写进整个集群Stream Load 表面上是把文件传到一台服务器实际走的是一条分布式写入路径FE 的http_port默认 8030接收请求后通过 HTTP 重定向把任务交给某个 BE 担任 Coordinator。Coordinator 按表的分布键把数据切分、散列到各 BE 并行落盘最后汇总执行结果返回给客户端。请求发给 FE 或某个 BE 都行发 FE 能借助轮询机制在集群内做负载均衡通常是更省心的选择。最小可用导入路径先用一张主键表接数据。主键表支持 upsert 语义——相同主键重复导入会覆盖旧值正好适合遥测数据同一设备反复上报的特点CREATE TABLE telemetry ( device_id BIGINT NOT NULL, metric VARCHAR(64) NOT NULL, reading DOUBLE, reading_time DATETIME NOT NULL ) ENGINEOLAP PRIMARY KEY(device_id, metric) DISTRIBUTED BY HASH(device_id);数据文件readings.csv长这样示例两行1001,cpu,82.5,2024-05-08 10:00:00 1002,mem,61.0,2024-05-08 10:00:05提交导入curl --location-trusted -u root: \ -H label:telemetry_20240508 \ -H column_separator:, \ -T readings.csv -XPUT \ http://fe_host:8030/api/demo/telemetry/_stream_loadURL 形态固定为/api/库名/表名/_stream_load-T指定本地文件-XPUT表示上传。导入完成会返回包含Status: Success、NumberLoadedRows、LoadTimeMs的 JSON此时数据已经可查。label是幂等键重复使用同一个 label 不会造成重复导入定时任务做失败重试时尤其有用。JSON 字段映射的三种姿势JSON 源文件的字段命名往往和表结构对不上配置就发生在这里。Stream Load 的映射有三档从轻到重同名直映JSON 的 key 与表列名一致时什么都不用配提取 改名jsonpaths声明从每个 JSON 对象里取哪些字段columns按顺序给它们起临时列名表达式转换在columns里直接写标量表达式落表前完成计算或类型转换。比如{device: 1001, temp: 25.3}要换算成华氏温度落库curl --location-trusted -u root: \ -H label:sensor_json \ -H format: json \ -H jsonpaths: [\$.device\, \$.temp\] \ -H columns: device_id, t, reading t * 9 / 5 32 \ -T sensors.json -XPUT \ http://fe_host:8030/api/demo/telemetry/_stream_load映射规则是前段按顺序、末段按列名jsonpaths与columns的临时列名一一对应最后再按名称落到表列上所以常见写法是前面几列直接改名最后一列做计算。完整的映射细则和返回值字段说明见 docs/en/loading/StreamLoad.md。合并提交调优参数速查高并发小批量单请求几 KB 到几十 MB是最容易翻车的形态每个请求都开一个事务、生成一个数据版本版本堆积会拖慢查询严重时触发too many versions报错。Merge Commitv3.4.0 起把时间窗内的多个同质请求合并成一个事务提交从源头上压住版本增速。curl --location-trusted -u root: \ -H enable_merge_commit:true \ -H merge_commit_interval_ms:5000 \ -H merge_commit_parallel:2 \ -T batch.csv -XPUT \ http://fe_host:8030/api/demo/telemetry/_stream_loadenable_merge_commit合并提交开关只对单表上的并发小批量场景有意义低并发下开启反而多等一个窗口merge_commit_interval_ms合并窗口时长窗口越大版本越少、可见延迟越高merge_commit_async异步模式下服务端收完数据立即返回不保证已落库需要自己后续查事务状态merge_commit_parallel合并窗口内的写入并行度。两个容易被忽略的细节能合并进同一事务的请求必须参数完全相同同质窗口内只要有一个请求数据不干净整批一起失败。合并模式下 label 由服务端生成你指定的会被忽略。故障排查速查表Stream Load 的返回 JSON 自带ErrorMsg和ErrorURL记录出错行明细结合下表多数失败一两分钟就能定位症状常见原因处理办法body exceed max size单文件超过 10 GB 上限streaming_load_max_mb拆分文件或调大 BE 参数需重启生效too many versions并发小批量导致版本增长过快启用 Merge Commit或降低导入频率导入频繁超时stream_load_default_timeout_second偏小默认 600 秒按数据量 ÷ 平均导入速度放宽单任务超时也可用timeout参数覆盖个别坏行拖垮整批默认严格模式不容忍错误行max_filter_ratio放宽到0.01左右别给太大否则真问题会被吞掉CSV 带表头导入报错老版本不会跳过首行v3.0 起加skip_header参数跳过导入之外查询加速Stream Load 的数据导入后即刻可查且落在表上的数据如果挂着物化视图视图会随导入同步更新实时看板可以直接搭在上面不用再垫一层 ETL。下一步建议先在测试集群把上面 curl 的完整流程走一遍重点核对返回 JSON 里每个字段的含义再往生产搬想进一步减少人工调度方向是把 Stream Load 接到 Routine Load 或任务调度系统里让导入变成无人值守的常驻链路。【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表