Vector Doris Sink 实战:通过 Stream Load 将日志批量写入 Apache Doris
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
Vector 的dorissink 负责将日志数据投递到 Apache Doris 数据库:它基于 Doris 的 Stream Load API,以“按 (database, table) 分区 → 攒批 → HTTP PUT 流式导入”的链路完成写入,并内置多端点负载均衡、健康检查、指数重试与 label 去重机制。读完本文,你可以掌握该 sink 的完整配置项、请求构造与重试/故障转移的实现细节,并能结合 src/sinks/doris/config.rs 等源码验证其运行行为。
组件定位与能力概览
Doris sink 的元数据定义在 website/cue/reference/components/sinks/doris.cue,其中声明的关键属性如下:
| 维度 | 取值 | 含义 |
|---|---|---|
delivery | at_least_once | 交付语义为“至少一次”,配合 label 机制尽量避免重复 |
development | beta | 组件仍处于 beta 阶段,行为可能随版本演进 |
egress_method | batch | 出口方式为批量发送,而非逐条推送 |
stateful | false | 无状态组件 |
acknowledgements | true | 支持端到端确认(E2E acks) |
healthcheck | enabled: true | 内置健康检查(基于 Doris bootstrap API) |
| 输入类型 | logs only | 仅接受 logs,不支持 metrics/traces |
在 src/sinks/doris/config.rs 中,SinkConfig实现将输入限定为Input::log(),与上述元数据一致。support.requirements明确要求 Doris 版本 1.0 或更高以获得最佳兼容性;website/cue/reference/services/doris.cue 中对 Doris 的定位是:一个现代 MPP 分析型数据库,提供亚秒级查询响应,适用于实时数仓、Ad-hoc 查询、统一数据分析与日志分析等 OLAP 场景。
工作原理:Stream Load 请求是如何构造的
从源码结构看,整个写入链路由以下模块协作完成(均位于 src/sinks/doris/):
- sink.rs:
DorisSink主循环,负责分区与攒批; - request_builder.rs:将批次编码、压缩后封装为
HttpRequest; - service.rs:
DorisService,真正发起 Stream Load 并解析响应; - client.rs:
DorisSinkClient,构造 HTTP 请求、处理重定向与健康检查; - health.rs、retry.rs:端点健康判定与重试策略;
- common.rs:多端点解析与公共配置。
事件分区与攒批
sink.rs中的DorisKeyPartitioner会把每个事件按渲染后的(database, table)元组(DorisPartitionKey)进行分区——这意味着database与table都是支持{{ }}模板的字段,一条流水线可以将不同事件路由到不同的库表。随后事件流经过batched_partitioned按batch配置(事件数、字节数、超时三者任一触发)成批,再交由DorisRequestBuilder编码(默认 newline-delimited + JSON,见下文)并构建 HTTP 请求。
Stream Load 请求细节
client.rs的build_request展示了请求的完整构造过程(client.rs#L121-L212):
- 请求路径为
PUT {scheme}://{authority}/api/{database}/{table}/_stream_load,其中 database/table 经过 RFC 3986 路径段百分号编码; - 固定设置
Content-Type: text/plain;charset=utf-8(Doris 对 JSON 数据期望纯文本 UTF-8)与Expect: 100-continue(Stream Load 协议要求先等待服务端确认再发送请求体); - 默认附带
label头,格式为{label_prefix}_{database}_{table}_{timestamp_ms}_{uuid}(见generate_label,client.rs#L94-L105)。Doris 用 label 识别并拒绝重复的 Stream Load 请求,从而在 Vector 重试时避免数据重复; - 若自定义
headers中包含group_commit: sync_mode或group_commit: async_mode,则不设置 label——group commit 会把多个导入合并进一个事务,label 在这种情况下没有意义(is_group_commit_enabled,client.rs#L112-L118); - 若配置了
compression,会追加Content-Encoding/Accept-Encoding头; - 认证信息(
auth)最后通过auth.apply附加到请求头。
重定向与响应判定
Doris Stream Load 的典型行为是 FE 节点将请求重定向到 BE 节点。send_stream_load(client.rs#L216-L335)实现了:
- 最多跟随 3 次重定向(301/302/307/308),并用
visited_urls集合检测重定向环,超限或成环时分别返回MaxRedirectsExceeded、RedirectLoop错误; - 最终响应体解析为 JSON,仅当
Status字段(不区分大小写)为success时判定为StreamLoadStatus::Successful,否则为Failure; - 该状态最终映射为事件终态:
Successful → EventStatus::Delivered,Failure → EventStatus::Errored(client.rs#L465-L472),供端到端确认机制使用。
健康检查
healthcheck_fenode通过GET {scheme}://{authority}/api/bootstrap探测 FE 节点(client.rs#L337-L418):HTTP 成功且响应 JSON 中msg == "success"才视为健康。端点级健康判定逻辑在 health.rs:成功状态码 → 健康;5xx → 不健康;其他状态码或网络错误则不做判定(None),交由上层持续观察。
完整配置参考
配置字段与类型的权威定义来自代码文档注释(config.rs#L27-L107)及其生成的 CUE 模式 website/cue/reference/components/sinks/generated/doris.cue。完整字段说明如下:
| 字段 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
endpoints | string 数组 | 是 | — | Doris 端点列表,必须带http/httpsscheme,可含主机与端口,如http://127.0.0.1:8030;配置多个端点即启用负载均衡 |
database | string(支持模板) | 是 | — | 目标数据库名,支持{{ }}模板,如mydatabase |
table | string(支持模板) | 是 | — | 目标表名,支持{{ }}模板,如mytable |
label_prefix | string | 否 | vector | Stream Load label 前缀,最终 label 为{label_prefix}_{database}_{table}_{timestamp}_{uuid} |
log_request | bool | 否 | false | 打开后以美化 JSON 记录每个 Stream Load 响应 |
headers | object(string→string) | 否 | {} | 自定义 HTTP 头,用于设置 Doris Stream Load 参数,如format(json/csv)、read_json_by_line、strip_outer_array、列映射等;也用于开启 group commit(group_commit: sync_mode/async_mode) |
encoding | object | 是(可省) | JSON + newline 分帧 | 编码配置;默认按事件序列化为 JSON 对象、多事件以换行分隔(NDJSON),见 config.rs#L126-L129 的Default实现 |
framing | object | 否 | newline delimited | 分帧配置 |
compression | enum | 否 | none | HTTP 请求压缩算法,可选none、gzip、snappy、zlib、zstd |
max_retries | int | 否 | -1 | 失败重试次数上限,-1表示无限重试 |
batch | object | 否 | 实时按字节默认值 | 事件批处理行为:max_events、max_bytes、timeout_secs任一达到即刷出;元数据声明的批量上限为 10 MB、超时 1.0s(doris.cue#L19-L23) |
auth | object | 否 | — | HTTP 认证策略(如user/password的 basic 认证);认证头在每次对 FE 的 Stream Load 请求中携带。注意:顶层auth与端点 URL 内嵌凭据不能同时使用,否则验证直接失败(见下文测试) |
request | object | 否 | 默认 | 出站请求的 Tower 中间件设置:并发、限流、超时与重试退避(退避策略遵循斐波那契数列) |
tls | object | 否 | — | TLS 配置,支持证书/主机名校验 |
distribution | object | 否 | — | Doris 端点健康判定选项(HealthConfig),控制多端点健康监控行为 |
acknowledgements | bool/object | 否 | false | 端到端确认行为控制 |
dangerously_allow_unconfined_template_resolution | bool | 否 | false | 危险选项:完全关闭模板 confinement 安全检查(见“模板安全”一节) |
一个典型的最小化配置(字段均可在上述源码中核实):
sinks: my_doris: type: doris inputs: [my_source] endpoints: - "http://127.0.0.1:8030" database: "mydatabase" table: "mytable" auth: user: "doris_user" password: "${DORIS_PASSWORD}" compression: "gzip" batch: max_events: 10000 timeout_secs: 1.0仓库中还生成了两份可直接参考的示例配置:minimal.yaml 与 advanced.yaml。
多端点:负载均衡与故障转移
how_it_works.load_balancing(doris.cue#L137-L151)描述了多端点语义:
- Round-robin 分发:请求在各可用 FE 端点间均匀分布;
- 健康监控:不健康端点自动被排除;
- 自动故障转移:某端点不可用时,流量切换到健康端点;
- 自动恢复:曾失败的端点会被周期性复测,恢复健康后重新纳入。
从源码结构看,这一行为由config.rs中build阶段的request_settings.distributed_service(DorisRetryLogic {}, services, health_config, DorisHealthLogic, ...)调用提供(config.rs#L281-L287):所有端点各自构建DorisService,再由统一的分布式服务包装层叠加重试逻辑(retry.rs)与健康调度。错误处理策略则包括:基于max_retries的自动重试(-1为无限)、端点间的重试退避以避免压垮集群、数据格式错误时记录错误并继续处理后续批次、以及网络层故障触发端点间切换。
配置校验:哪些配置会被拒绝
ValidatedSink::validate(config.rs#L175-L217)在组件启动前执行一系列检查,失败即阻止流水线构建:
endpoints不能为空;- 每个端点 scheme 必须是
http或https,且必须包含主机名; - 顶层
auth与端点 URL 中内嵌的user:pass@凭据不允许同时出现(auth.choose_one); database/table模板需通过 confinement 检查(见下节)。
这些规则均有对应的单元测试守护:validate_rejects_non_http_scheme验证非 http scheme(如ftp://)被拒绝;validate_rejects_auth_conflict_with_endpoint_credentials验证凭据冲突报错;test_default_values验证label_prefix默认vector、max_retries默认-1(config.rs#L329-L406)。此外,validate刻意只做纯端点检查,把需要读取证书磁盘文件的DorisCommon解析推迟到build阶段,从而保持vector validate --no-environment无文件系统依赖(见 config.rs#L180-L182 注释)。
模板安全(Confinement)
由于database/table支持模板,恶意日志事件理论上可以改写写入目标。Vector 通过模板 confinement 机制设防:没有可用静态前缀的裸模板(如{{ tenant }})在验证期即被拒绝;带前缀的模板(如mydb_{{ tenant }})在渲染期还会阻止../这类路径逃逸。对应测试见 config.rs#L408-L443:confinement_rejects_unconfined_database_template、confinement_blocks_dotdot_escape_at_render。若确需无约束模板,可显式设置dangerously_allow_unconfined_template_resolution: true——该选项会绕过启动期与运行期所有 confinement 检查,意味着能控制模板字段的日志生产者可以写入任意库表,生产环境应慎用。
可观测性:内部遥测事件
DorisService::reporter_run(service.rs#L36-L82)在收到成功响应后会从 Stream Load 结果 JSON 中提取并上报内部指标:
DorisRowsLoaded:携带LoadBytes与NumberLoadedRows,反映本批实际加载的字节数与行数;DorisRowsFiltered:当NumberFilteredRows > 0时上报被 Doris 过滤(拒绝/丢弃)的行数,方便及时发现数据质量问题;- 开启
log_request后,每次响应会以格式化 JSON 记录在 INFO 日志中(含 HTTP 状态码、Stream Load 状态与完整响应体)。
此外send_stream_load每成功送达一批还会发出EndpointBytesSent事件,记录端点、协议(http/https)与字节数(client.rs#L299-L303)。
验证与深入测试
- 配置层单测:
cargo test --lib doris可运行 config.rs 中的validate_produces_usable_values、凭据冲突与非 http scheme 拒绝等测试; - 集成测试:src/sinks/doris/integration_test.rs 提供贴近真实 Doris 服务行为的集成用例;
- 端到端确认:由于组件支持
acknowledgements,可将其放入启用 ack 的流水线中验证事件终态与EventStatus::Delivered/Errored的映射(service.rs#L128-L133)。
需要注意的适用前提:该 sink 为 beta 组件、仅支持 logs 输入、交付语义为“至少一次”(label 机制用于降低重试导致的重复,但从组件元数据的分类看,官方声明仍为 at_least_once);group commit 开启后 label 会被跳过,去重保护也随之失效,这是以可见性/合并导入换取吞吐的取舍,选型时应结合 Doris 集群能力权衡。
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考