尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

从 Elasticsearch 迁移到 Apache Doris:日志分析与可观测性平台实战

从 Elasticsearch 迁移到 Apache Doris:日志分析与可观测性平台实战 1. 为什么我把日志系统从 Elasticsearch 换成了 Apache Doris先交代背景。我负责维护一套中等规模的微服务集群大概四十多个服务实例日均日志量在 800GB 到 1.2TB 之间浮动。最早用的是经典的 ELK 组合Elasticsearch 存日志、Logstash 做采集、Kibana 做展示。这套方案用了两年多问题越来越明显磁盘占用高得离谱同样的原始日志量ES 的索引膨胀能到原始数据的 2.5 到 3 倍查询稍微复杂一点比如做多维度聚合或者跨天的大范围扫描响应时间直接飙到十几秒甚至几十秒集群节点一扩再扩成本压不住。后来我开始认真评估 Apache Doris。说实话一开始我是有顾虑的Doris 在 OLAP 场景的口碑一直不错但拿它当日志存储和可观测性底座当时心里没底。真正让我下决心的是做了一轮压测同样的日志数据Doris 的存储压缩比能做到 5 到 8 倍复杂聚合查询基本在亚秒到两三秒之间返回。这个差距不是一点半点。这篇笔记就是把这套迁移和落地的完整过程拆开讲清楚。Apache Doris 做日志分析与可观测性核心思路是用一套 MPP 架构的分析型数据库同时扛住日志存储、全文检索、指标聚合和链路追踪这几件事而不是像传统方案那样 ES 管日志、Prometheus 管指标、Jaeger 管链路各搞一套。适合谁看如果你正在被日志系统的成本和查询性能折磨或者你正在设计一套新的可观测性平台这篇内容应该能帮你少走一些弯路。我会把架构选型的逻辑、表结构设计、数据写入链路、查询优化、踩过的坑都摊开讲。2. 整体架构设计与选型逻辑拆解2.1 为什么是 Doris 而不是继续堆 ES 节点先把这个最核心的问题说清楚。ES 的问题不在于它不好而在于它的设计目标和日志分析场景之间存在错位。ES 底层是 Lucene倒排索引是它的核心这决定了它在全文检索上很强但倒排索引本身非常吃存储。日志数据有个特点写多读少、冷数据占比极高、查询模式以聚合统计和时间范围过滤为主。你拿一个为全文检索优化的引擎去干聚合分析的活就像开越野车送快递能送但油耗和维护成本都不划算。Doris 的定位是 MPP 分析型数据库列式存储加向量化执行引擎天然适合大宽表的聚合扫描。它的数据模型里有几个关键设计对日志场景特别友好列式存储 高压缩比日志字段里大量重复值比如服务名、日志级别、主机名列存下同列数据连续存放压缩算法能发挥到极致。我实测 ZSTD 压缩下原始 JSON 日志 1TB 压到 130GB 左右。分区 分桶两级切分按天分区、按服务名或主机名分桶查询时分区裁剪加桶裁剪扫描的数据量能砍掉一大截。物化视图和 Rollup可以针对常用聚合查询预计算比如按分钟统计各服务 ERROR 数量直接命中 Rollup 表不用扫原始数据。倒排索引支持Doris 从 2.0 开始支持倒排索引虽然不如 ES 那么成熟但应付日志里的关键字检索足够了。注意Doris 的倒排索引和 ES 的倒排索引不是一回事。Doris 的倒排索引是辅助加速手段不是存储的核心结构所以不会带来 ES 那种存储膨胀问题。2.2 可观测性三大支柱怎么统一到一套存储可观测性通常讲三大支柱Logs日志、Metrics指标、Traces链路。传统做法是三套系统各管一摊查询的时候要在三个界面之间来回跳关联分析基本靠人肉。用 Doris 的思路是三者的数据模型虽然不同但本质上都是带时间戳的结构化或半结构化数据完全可以统一到一套存储里用不同的表来承载。我的做法是建三张核心表表名承载数据数据模型分区策略分桶策略logs_detail原始日志明细Duplicate 模型按天分区按 service_name 哈希分桶metrics_rollup聚合指标Aggregate 模型按天分区按 metric_name 哈希分桶traces_span链路 Span 数据Duplicate 模型按天分区按 trace_id 哈希分桶这样做的好处是做关联查询的时候不用跨系统。比如我想查某个 trace_id 对应的所有日志和它经过的每个服务的耗时指标一条 SQL 就能搞定不用先去 Jaeger 查链路再去 Kibana 查日志。2.3 写入链路的设计取舍写入链路我选的是Flink Doris Stream Load的组合。日志从各个服务通过 Filebeat 采集到 KafkaFlink 消费 Kafka 做清洗和字段提取然后通过 Stream Load 批量写入 Doris。为什么不用 Routine Load 直接消费 KafkaRoutine Load 确实更简单Doris 原生支持不用额外维护 Flink 作业。但我的场景里有几个需求 Routine Load 满足不了一是需要在写入前做字段解析和脱敏比如把 JSON 日志里的嵌套字段拍平、把手机号邮箱做掩码二是需要做多流合并把日志流和链路流按 trace_id 做关联后写入三是需要动态路由不同服务的日志写到不同的分区。这些用 Flink 做更灵活。写入批次的大小很关键。我试过几个档位批次太小比如 1000 条一批Stream Load 的导入频率太高Doris 的 compaction 压力大查询性能会抖动批次太大比如 10 万条一批延迟又上去了日志从产生到可查要等好几分钟。最后定在每批 2 万到 5 万条或者每 10 秒触发一次兼顾延迟和写入效率。3. 核心表结构设计与实操要点3.1 日志明细表字段类型和索引怎么定日志表的设计直接决定了后面查询爽不爽。我踩过的最大坑是字段类型选得太随意导致后面查询要么慢要么报错。下面是我最终定下来的建表语句核心部分CREATE TABLE logs_detail ( ts DATETIMEV2(3) NOT NULL COMMENT 日志时间戳毫秒精度, service_name VARCHAR(64) NOT NULL COMMENT 服务名, host_ip VARCHAR(32) NOT NULL COMMENT 主机IP, log_level VARCHAR(8) NOT NULL COMMENT 日志级别, trace_id VARCHAR(64) NULL COMMENT 链路追踪ID, span_id VARCHAR(32) NULL COMMENT Span ID, message TEXT NULL COMMENT 日志正文, extra JSON NULL COMMENT 扩展字段, INDEX idx_message (message) USING INVERTED COMMENT 正文倒排索引, INDEX idx_trace (trace_id) USING INVERTED COMMENT trace_id 倒排索引 ) ENGINEOLAP DUPLICATE KEY(ts, service_name, host_ip) PARTITION BY RANGE(ts) () DISTRIBUTED BY HASH(service_name) BUCKETS 32 PROPERTIES ( compression ZSTD, replication_num 3, dynamic_partition.enable true, dynamic_partition.time_unit DAY, dynamic_partition.start -30, dynamic_partition.end 3, dynamic_partition.prefix p, dynamic_partition.buckets 32 );几个关键决策点展开说时间戳用 DATETIMEV2(3) 而不是 BIGINT。用 BIGINT 存毫秒时间戳确实省空间但查询的时候每次都要写from_unixtime(ts/1000)非常别扭而且分区裁剪用不上。DATETIMEV2 支持毫秒精度直接可读分区裁剪也正常。存储上多占的那点空间相比查询便利性完全值得。message 字段用 TEXT 而不是 VARCHAR。日志正文长度不可控有的堆栈信息能到几十KB。VARCHAR 有长度上限超了会截断或者报错。TEXT 没有这个问题配合倒排索引查询效率也可以接受。分桶数选 32 是基于数据量算出来的。Doris 官方建议单个桶的数据量在 1GB 到 10GB 之间。我日均 1TB 原始数据压缩后约 150GB按 30 天保留算单天分区压缩后约 5GB。32 个桶每个桶约 160MB偏小但问题不大因为查询并发高的时候桶多一点能更好利用多核。如果你数据量更大可以按这个公式算分桶数 单分区压缩后大小 / 期望单桶大小。倒排索引只建在 message 和 trace_id 上。倒排索引不是免费的每个索引都会增加写入开销和存储占用。log_level、service_name 这种低基数字段用前缀索引或者 Bloom Filter 就够了没必要上倒排。3.2 指标聚合表Aggregate 模型怎么用指标数据的特点是写入频繁、查询以聚合为主、历史数据很少查明细。这种场景用 Aggregate 模型最合适Doris 会在后台自动做预聚合查询的时候直接读聚合结果。CREATE TABLE metrics_rollup ( ts DATETIMEV2(3) NOT NULL, metric_name VARCHAR(128) NOT NULL, service_name VARCHAR(64) NOT NULL, labels VARCHAR(512) NULL, value_sum DOUBLE SUM DEFAULT 0, value_count BIGINT SUM DEFAULT 0, value_max DOUBLE MAX DEFAULT 0, value_min DOUBLE MIN DEFAULT 0 ) ENGINEOLAP AGGREGATE KEY(ts, metric_name, service_name, labels) PARTITION BY RANGE(ts) () DISTRIBUTED BY HASH(metric_name) BUCKETS 16 PROPERTIES ( compression ZSTD, replication_num 3, dynamic_partition.enable true, dynamic_partition.time_unit DAY, dynamic_partition.start -90, dynamic_partition.end 3, dynamic_partition.prefix p );这里有个设计上的取舍聚合粒度定在多少。定得太细比如秒级预聚合效果不明显存储也大定得太粗比如小时级查询灵活性差。我最后定在分钟级写入的时候 Flink 做一层分钟级预聚合再写入Doris 的 Aggregate 模型再做一层合并。这样既保证了查询能下钻到分钟又控制了数据量。labels字段我用的是拼接后的字符串比如envprod,regioncn-east,instance10.0.1.5:9090。为什么不拆成多个列因为 Prometheus 的 label 是动态的不同指标 label 集合不一样拆列的话表结构没法固定。拼成字符串后用LIKE或者split_by_string函数过滤配合倒排索引效率可以接受。3.3 链路表trace_id 和 span 的存储优化链路数据的查询模式很特殊要么按 trace_id 查整条链路的所有 span要么按时间范围加服务名查慢请求。前者要求 trace_id 的查询极快后者要求时间范围扫描高效。CREATE TABLE traces_span ( ts DATETIMEV2(3) NOT NULL, trace_id VARCHAR(64) NOT NULL, span_id VARCHAR(32) NOT NULL, parent_span_id VARCHAR(32) NULL, service_name VARCHAR(64) NOT NULL, operation_name VARCHAR(256) NULL, duration_ms BIGINT NOT NULL, status_code INT NOT NULL, tags JSON NULL, INDEX idx_trace (trace_id) USING INVERTED ) ENGINEOLAP DUPLICATE KEY(ts, trace_id, span_id) PARTITION BY RANGE(ts) () DISTRIBUTED BY HASH(trace_id) BUCKETS 32 PROPERTIES ( compression ZSTD, replication_num 3, dynamic_partition.enable true, dynamic_partition.time_unit DAY, dynamic_partition.start -15, dynamic_partition.end 3, dynamic_partition.prefix p );链路数据我只保留 15 天比日志的 30 天短。原因是链路数据量虽然比日志小但查询频率也低大部分时候只有排查问题才会查15 天足够覆盖绝大多数排查场景。如果你有合规要求需要保留更久可以调大dynamic_partition.start的绝对值。分桶键选trace_id而不是service_name是因为链路查询最核心的场景是按 trace_id 查全链路按 trace_id 分桶能让同一个 trace 的所有 span 落在同一个桶里查询时只扫一个桶。4. 数据写入链路与查询优化的实战细节4.1 Flink 写入作业的关键配置Flink 作业这块我调了不少参数下面这几个是影响最大的// Doris Stream Load 的 Sink 配置 DorisSink.Builder builder DorisSink.builder() .setDorisReadOptions(DorisReadOptions.builder().build()) .setDorisExecutionOptions(DorisExecutionOptions.builder() .setLabelPrefix(logs_ System.currentTimeMillis()) .setDeletable(false) .setBufferCount(5) // 攒够5个批次再提交 .setBufferSize(64 * 1024 * 1024) // 单批次最大64MB .setMaxRetries(3) .setStreamLoadProp(props) .build()) .setDorisOptions(DorisOptions.builder() .setFenodes(fe1:8030,fe2:8030,fe3:8030) .setTableIdentifier(ods.logs_detail) .setUsername(log_writer) .setPassword(******) .build());bufferCount和bufferSize这两个参数要配合调。我一开始只设了bufferSize没设bufferCount结果小批次频繁提交Doris 的 compaction 跟不上查询延迟从 1 秒涨到 8 秒。后来加上bufferCount5让 Flink 攒够 5 个批次再一次性提交compaction 压力明显下降。还有一个容易忽略的点Stream Load 的 label 要保证唯一。Doris 用 label 做导入事务的去重如果 label 重复导入会被拒绝。我用logs_ 时间戳 随机数的格式生成 label基本不会冲突。4.2 查询优化的几个实用手段分区裁剪是第一优先级。所有查询必须带时间范围条件而且时间字段要直接用分区列不能包函数。比如WHERE ts 2024-01-01 00:00:00能触发分区裁剪WHERE date(ts) 2024-01-01就不行因为分区裁剪发生在执行计划生成阶段函数包裹后优化器识别不了。物化视图预计算高频聚合。我建了几个物化视图比如按分钟统计各服务的 ERROR 日志数量CREATE MATERIALIZED VIEW mv_error_count_by_minute REFRESH ASYNC EVERY(INTERVAL 1 MINUTE) DISTRIBUTED BY HASH(service_name) BUCKETS 8 PROPERTIES (replication_num 3) AS SELECT date_trunc(ts, minute) AS minute_ts, service_name, count(*) AS error_count FROM logs_detail WHERE log_level ERROR GROUP BY minute_ts, service_name;这个物化视图建了之后告警系统查 ERROR 数量的查询从平均 2.3 秒降到 80 毫秒。代价是每分钟有一次后台刷新任务会消耗一些资源但相比查询性能的提升完全值得。倒排索引的查询写法有讲究。Doris 的倒排索引对MATCH_ANY、MATCH_ALL、MATCH_PHRASE这几个函数支持最好。比如查 message 里包含 timeout 或 connection refused 的日志SELECT ts, service_name, message FROM logs_detail WHERE ts 2024-01-01 00:00:00 AND ts 2024-01-01 01:00:00 AND message MATCH_ANY timeout connection refused LIMIT 100;注意MATCH_ANY后面跟的是空格分隔的关键词列表不是 SQL 的 OR 语法。这个写法能命中倒排索引比message LIKE %timeout% OR message LIKE %connection refused%快一个数量级。4.3 冷热数据分层存储日志数据有个特点最近 3 天的查询量占 90% 以上7 天前的数据基本没人查。针对这个特点我做了冷热分层热数据最近 3 天存在 SSD 上副本数 3保证查询性能和可用性。温数据4 到 15 天存在 HDD 上副本数 2查询性能可以接受。冷数据16 到 30 天存到对象存储上用 Doris 的冷热分层功能自动下沉副本数 1。Doris 的冷热分层通过storage_policy配置ALTER TABLE logs_detail SET ( storage_policy hot_to_cold_policy );对应的 policy 定义CREATE STORAGE POLICY hot_to_cold_policy PROPERTIES( storage_resource remote_oss, cooldown_ttl 259200 -- 3天后下沉到冷存储 );这套分层下来存储成本比全 SSD 方案降了大概 60%而热数据的查询性能没有受影响。5. 常见问题与排查技巧实录5.1 写入延迟突然飙升怎么排查这是我最常遇到的问题表现是日志从产生到可查的时间从正常的 10 秒涨到几分钟。排查思路按这个顺序走第一步看 Flink 作业的反压情况。如果 Flink 的 Sink 算子出现反压说明写入 Doris 的速度跟不上消费 Kafka 的速度。这时候去 Doris 的 FE 看 Stream Load 的导入任务状态大概率能看到大量LOADING状态的任务堆积。第二步看 Doris 的 compaction 分数。执行SHOW TABLET FROM logs_detail看MaxCompactionScore字段如果超过 100 说明 compaction 严重滞后。compaction 滞后的原因是小文件太多通常是因为 Stream Load 批次太小或者频率太高。第三步看 BE 的磁盘 IO。如果磁盘 IO 打满写入自然慢。这时候要么加磁盘要么调整写入策略减少 IO 压力。我遇到过一次比较典型的情况某个服务突然开始疯狂打日志日志量涨了 10 倍Flink 作业的写入批次大小没变导致 Stream Load 频率暴涨Doris 的 compaction 直接崩了。解决办法是给 Flink 作业加了一个动态限流逻辑当检测到 Kafka 消费延迟超过阈值时自动增大批次大小、降低提交频率。5.2 查询变慢的几种典型原因查询变慢的原因比写入延迟更分散我整理了一个速查表现象可能原因排查方法解决手段简单查询也慢分区裁剪失效EXPLAIN看扫描分区数检查 WHERE 条件是否直接用了分区列聚合查询慢没命中物化视图EXPLAIN看是否命中 Rollup调整物化视图定义或查询写法全文检索慢倒排索引未命中EXPLAIN看是否走倒排索引改用 MATCH_ANY 等函数偶发性慢查询资源竞争看 FE 的查询队列和 BE 的负载加资源组隔离或错峰查询扫描数据量大分桶数不合理看 Tablet 的扫描行数调整分桶数或加分区提示EXPLAIN是排查 Doris 查询问题的第一工具任何查询优化之前先看执行计划不要凭感觉调。5.3 几个我踩过的坑坑一DATETIME 精度不够导致日志乱序。最早我用的是 DATETIME秒级精度结果同一秒内的多条日志排序不稳定排查问题时顺序对不上。后来改成 DATETIMEV2(3) 毫秒精度才解决。日志场景建议至少毫秒精度如果日志量特别大且对顺序敏感可以考虑微秒精度。坑二JSON 字段查询性能差。我一开始把扩展字段都塞到 JSON 列里查询的时候用extra[user_id]这种写法。后来发现 JSON 字段的查询没法走索引每次都要全表扫。解决办法是把高频查询的 JSON 字段提取成独立的列JSON 列只存低频查询的字段。坑三动态分区创建失败导致写入中断。动态分区依赖 FE 的调度如果 FE 负载高或者调度延迟可能出现分区还没创建但数据已经写入的情况导致导入失败。我的解决办法是在 Flink 作业里加了重试逻辑遇到分区不存在的错误时等待几秒重试同时给 FE 的动态分区调度加了监控告警。坑四副本数设太高导致写入放大。我一开始所有表都设了replication_num3写入的时候每个副本都要写一遍写入放大 3 倍。后来把冷数据和温数据的副本数降到 2 和 1写入压力明显下降。副本数的选择要在可用性和写入成本之间权衡不是越高越好。6. 这套方案的实际收益与扩展思路6.1 收益的量化对比迁移完成稳定运行三个月后我做了一次全面的对比指标原 ELK 方案Doris 方案变化存储占用30天约 75TB约 12TB降低 84%日均写入延迟 P9945 秒8 秒降低 82%简单聚合查询 P953.2 秒0.4 秒降低 87%复杂聚合查询 P9518 秒2.1 秒降低 88%全文检索 P951.1 秒0.9 秒基本持平集群节点数12 节点6 节点减少 50%月度成本基准约 35%降低 65%全文检索这块 Doris 确实不如 ES但差距没有想象中大而且日志场景里真正的全文检索需求占比很低大部分查询都是带条件的结构化过滤加聚合。用 35% 的成本换来其他指标的全面碾压这笔账怎么算都划算。6.2 后续可以扩展的方向这套架构跑稳之后我又陆续加了一些扩展告警规则引擎。基于 Doris 的物化视图和定时查询做了一个轻量级的告警引擎。比如每分钟查一次 ERROR 日志数量超过阈值就触发告警。相比 Prometheus 的告警规则用 SQL 写告警规则更灵活能做更复杂的条件组合。日志聚类分析。用 Doris 的simhash或者minhash函数对日志正文做相似度计算把相似的日志聚成一类帮助快速发现异常模式。这个功能在排查突发故障时特别有用能从海量日志里快速定位到异常的那几类。与 Grafana 集成。Doris 有官方的 Grafana 数据源插件可以直接把 Doris 作为 Grafana 的数据源用 SQL 查出来的结果直接画图。这样指标、日志、链路都能在 Grafana 里统一展示不用再维护 Kibana 那套。数据归档到对象存储。超过 30 天的日志通过 Doris 的EXPORT功能导出到对象存储上的 Parquet 文件需要的时候再用Broker Load导回来。这样既满足了合规保留要求又不占用 Doris 的存储资源。6.3 一些个人体会这套方案不是银弹它适合的是日志量中等偏上、查询以聚合分析为主、对成本敏感的场景。如果你的日志量很小日均几十GB或者全文检索需求占比很高那 ES 或者专门的日志服务可能更合适。技术选型永远要看具体场景不要因为 Doris 现在火就无脑上。另外Doris 的社区版本迭代很快我用的 2.1 版本相比 2.0 在倒排索引和冷热分层上都有明显改进。建议保持关注版本更新但生产环境升级要谨慎先在测试环境跑一轮完整的回归测试再上。最后说一个我觉得最重要的经验可观测性平台的价值不在于技术多先进而在于能不能让排查问题的人快速拿到想要的数据。我见过太多团队花大价钱搭了一套很炫的平台结果开发排查问题时还是习惯直接上机器grep日志。工具再好没人用就是白搭。所以在设计这套系统的时候我花了很多时间在查询模板和自助分析上让开发人员不用写复杂的 SQL 就能查到想要的日志。这个投入的回报比任何技术优化都高。
返回列表