ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Vector Doris Sink 实战:通过 Stream Load 将日志批量写入 Apache Doris

Vector Doris Sink 实战:通过 Stream Load 将日志批量写入 Apache Doris Vector Doris Sink 实战通过 Stream Load 将日志批量写入 Apache Doris【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vectorVector 的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其中声明的关键属性如下维度取值含义deliveryat_least_once交付语义为“至少一次”配合 label 机制尽量避免重复developmentbeta组件仍处于 beta 阶段行为可能随版本演进egress_methodbatch出口方式为批量发送而非逐条推送statefulfalse无状态组件acknowledgementstrue支持端到端确认E2E ackshealthcheckenabled: 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.rsDorisSink主循环负责分区与攒批request_builder.rs将批次编码、压缩后封装为HttpRequestservice.rsDorisService真正发起 Stream Load 并解析响应client.rsDorisSinkClient构造 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;charsetutf-8Doris 对 JSON 数据期望纯文本 UTF-8与Expect: 100-continueStream Load 协议要求先等待服务端确认再发送请求体默认附带label头格式为{label_prefix}_{database}_{table}_{timestamp_ms}_{uuid}见generate_labelclient.rs#L94-L105。Doris 用 label 识别并拒绝重复的 Stream Load 请求从而在 Vector 重试时避免数据重复若自定义headers中包含group_commit: sync_mode或group_commit: async_mode则不设置 label——group commit 会把多个导入合并进一个事务label 在这种情况下没有意义is_group_commit_enabledclient.rs#L112-L118若配置了compression会追加Content-Encoding/Accept-Encoding头认证信息auth最后通过auth.apply附加到请求头。重定向与响应判定Doris Stream Load 的典型行为是 FE 节点将请求重定向到 BE 节点。send_stream_loadclient.rs#L216-L335实现了最多跟随 3 次重定向301/302/307/308并用visited_urls集合检测重定向环超限或成环时分别返回MaxRedirectsExceeded、RedirectLoop错误最终响应体解析为 JSON仅当Status字段不区分大小写为success时判定为StreamLoadStatus::Successful否则为Failure该状态最终映射为事件终态Successful → EventStatus::DeliveredFailure → EventStatus::Erroredclient.rs#L465-L472供端到端确认机制使用。健康检查healthcheck_fenode通过GET {scheme}://{authority}/api/bootstrap探测 FE 节点client.rs#L337-L418HTTP 成功且响应 JSON 中msg success才视为健康。端点级健康判定逻辑在 health.rs成功状态码 → 健康5xx → 不健康其他状态码或网络错误则不做判定None交由上层持续观察。完整配置参考配置字段与类型的权威定义来自代码文档注释config.rs#L27-L107及其生成的 CUE 模式 website/cue/reference/components/sinks/generated/doris.cue。完整字段说明如下字段类型必填默认值说明endpointsstring 数组是—Doris 端点列表必须带http/httpsscheme可含主机与端口如http://127.0.0.1:8030配置多个端点即启用负载均衡databasestring支持模板是—目标数据库名支持{{ }}模板如mydatabasetablestring支持模板是—目标表名支持{{ }}模板如mytablelabel_prefixstring否vectorStream Load label 前缀最终 label 为{label_prefix}_{database}_{table}_{timestamp}_{uuid}log_requestbool否false打开后以美化 JSON 记录每个 Stream Load 响应headersobjectstring→string否{}自定义 HTTP 头用于设置 Doris Stream Load 参数如formatjson/csv、read_json_by_line、strip_outer_array、列映射等也用于开启 group commitgroup_commit: sync_mode/async_modeencodingobject是可省JSON newline 分帧编码配置默认按事件序列化为 JSON 对象、多事件以换行分隔NDJSON见 config.rs#L126-L129 的Default实现framingobject否newline delimited分帧配置compressionenum否noneHTTP 请求压缩算法可选none、gzip、snappy、zlib、zstdmax_retriesint否-1失败重试次数上限-1表示无限重试batchobject否实时按字节默认值事件批处理行为max_events、max_bytes、timeout_secs任一达到即刷出元数据声明的批量上限为 10 MB、超时 1.0sdoris.cue#L19-L23authobject否—HTTP 认证策略如user/password的 basic 认证认证头在每次对 FE 的 Stream Load 请求中携带。注意顶层auth与端点 URL 内嵌凭据不能同时使用否则验证直接失败见下文测试requestobject否默认出站请求的 Tower 中间件设置并发、限流、超时与重试退避退避策略遵循斐波那契数列tlsobject否—TLS 配置支持证书/主机名校验distributionobject否—Doris 端点健康判定选项HealthConfig控制多端点健康监控行为acknowledgementsbool/object否false端到端确认行为控制dangerously_allow_unconfined_template_resolutionbool否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_balancingdoris.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::validateconfig.rs#L175-L217在组件启动前执行一系列检查失败即阻止流水线构建endpoints不能为空每个端点 scheme 必须是http或https且必须包含主机名顶层auth与端点 URL 中内嵌的user:pass凭据不允许同时出现auth.choose_onedatabase/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默认-1config.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-L443confinement_rejects_unconfined_database_template、confinement_blocks_dotdot_escape_at_render。若确需无约束模板可显式设置dangerously_allow_unconfined_template_resolution: true——该选项会绕过启动期与运行期所有 confinement 检查意味着能控制模板字段的日志生产者可以写入任意库表生产环境应慎用。可观测性内部遥测事件DorisService::reporter_runservice.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_oncegroup commit 开启后 label 会被跳过去重保护也随之失效这是以可见性/合并导入换取吞吐的取舍选型时应结合 Doris 集群能力权衡。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进