ARTICLE · INTELLIGENCE

战地情报 · 详情页

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

Apache Doris实战:构建海量时空数据分析平台的全链路方案

Apache Doris实战:构建海量时空数据分析平台的全链路方案 最近在开发一个海洋环境监测系统时遇到了一个棘手的问题如何高效地处理和分析来自不同传感器、格式各异的“海”量时空数据。这些数据不仅体量庞大而且维度复杂传统的数据库和简单的文件存储方案在查询效率和扩展性上都遇到了瓶颈。经过一番技术选型和实践我最终选择将数据导入Apache Doris进行分析其卓越的OLAP性能和易用性彻底解决了我们的痛点。本文将完整分享这套从原始数据到可视化分析的全链路实战方案涵盖数据模型设计、多种数据导入方式、实时聚合查询以及性能调优要点无论你是数据分析师、后端开发还是架构师都能从中获得可直接复用于项目的经验。1. 背景与核心概念为什么选择 Doris 处理“海”量数据在物联网、互联网、金融等领域我们每天都在产生“海”量数据。这里的“海”并不仅仅指数据体积大TB/PB级更意味着数据产生的速度极快流式、结构多样结构化、半结构化且价值密度不均。处理这类数据传统的事务型数据库如MySQL在复杂分析查询上力不从心而早期的大数据方案如Hadoop生态又往往架构复杂、运维成本高。Apache Doris是一个基于 MPP 架构的高性能、实时的分析型数据库。它主要解决了海量数据的在线分析处理OLAP问题。其核心优势在于极速查询即使面对亿级甚至十亿级数据大部分聚合查询也能在亚秒级返回结果。兼容MySQL协议使用标准SQL语法并兼容MySQL通信协议降低了学习和使用门槛现有BI工具如FineBI、Tableau和应用程序可以无缝接入。简化架构一个系统同时支持高吞吐的批量数据导入和低延迟的实时数据导入无需维护复杂的 Lambda 或 Kappa 架构。易运维支持在线弹性扩缩容并具备完善的监控体系。简单来说当你的业务面临“数据量大、查询复杂、要求响应快”的挑战时Doris 是一个非常值得考虑的解决方案。接下来我们将从零开始搭建一个完整的海洋监测数据分析平台。2. 环境准备与版本说明在开始实战之前需要准备好运行环境。本文示例基于以下环境但核心步骤和原理适用于其他版本。操作系统CentOS 7.9 或 Ubuntu 20.04 LTS建议使用Linux服务器Apache Doris2.0.3当前较新的稳定版本JavaJDK 11Doris FE/BE 依赖示例数据模拟的海洋传感器数据CSV格式客户端工具mysql-client用于通过MySQL协议连接Doris。DBeaver/DataGrip图形化数据库管理工具可选。版本兼容性说明Doris 2.x 版本在数据导入、查询优化和生态集成上相比 1.x 有显著提升。部署时请务必参考 Apache Doris 官网 的发布说明确认组件间的版本依赖。生产环境建议使用奇数版本如2.0.x的次新版本以平衡新特性与稳定性。3. 核心原理与数据模型设计在导入数据前合理的数据模型设计是发挥 Doris 性能的关键。Doris 主要支持两种数据模型Duplicate Key 模型和Aggregate Key 模型包括 Unique Key 模型。此外分区和分桶是影响数据分布和查询效率的核心机制。3.1 数据模型选择假设我们的海洋传感器数据包含以下字段sensor_id(传感器编号)timestamp(数据采集时间戳)temperature(水温)salinity(盐度)ph(酸碱度)location(经纬度如POINT(120.5 30.3))场景分析需要存储最细粒度原始数据并可能基于任意字段进行过滤查询。 - 选择Duplicate Key 模型。它不对数据做任何聚合保留完整的导入数据行。需要按维度实时聚合例如实时查看每个传感器的最新状态或每小时的平均温度。 - 选择Aggregate Key 模型。它会在数据导入时根据 Key 列自动进行预聚合如 SUM, MAX, MIN, REPLACE。对于监测场景我们通常需要两种能力查询原始明细和查看聚合报表。因此可以创建两张表或使用 Aggregate 模型中的REPLACE聚合方式保存最新状态。3.2 分区与分桶分区Partitioning常用于按时间范围如天、月划分数据。这可以实现分区裁剪查询时只扫描相关分区极大提升性能。对于时序数据按timestamp的日期dt分区是标准做法。分桶Bucketing在分区内数据被进一步划分为多个 Tablet数据分片。分桶列的选择对查询性能至关重要应选择高频查询条件或 Join 条件的列如sensor_id。分桶数建议为机器磁盘数量的整数倍单个 Tablet 数据量在 100MB-1GB 为宜。4. 完整实战构建海洋监测数据分析平台4.1 部署 Apache Doris 集群单机伪集群示例首先从官网下载 Doris 安装包并解压。Doris 包含 FEFrontend和 BEBackend两种角色。1. 启动 FE元数据管理与查询协调# 进入FE目录 cd fe # 修改配置文件 conf/fe.conf指定元数据目录和JAVA_HOME如果未全局设置 # 初始化FE元数据 ./bin/start_fe.sh --daemon使用 MySQL 客户端连接 FE默认端口 9030并设置 root 密码mysql -h 127.0.0.1 -P 9030 -uroot # 在MySQL客户端内执行 SET PASSWORD FOR root PASSWORD(your_password);2. 启动 BE数据存储与计算# 进入BE目录 cd be # 修改配置文件 conf/be.conf主要配置 storage_root_path数据存储路径 ./bin/start_be.sh --daemon3. 添加 BE 节点到集群再次连接 FE执行以下 SQLALTER SYSTEM ADD BACKEND “你的服务器IP:9050“;通过SHOW BACKENDS\G命令检查 BE 状态是否正常。4.2 创建数据库与数据表连接 Doris 后我们为海洋监测数据创建数据库和表。-- 创建数据库 CREATE DATABASE IF NOT EXISTS ocean_monitor; USE ocean_monitor; -- 创建一张明细表Duplicate Key 模型用于存储原始传感器数据。 -- 我们按天分区按传感器ID分桶。 CREATE TABLE IF NOT EXISTS sensor_data_detail ( sensor_id INT NOT NULL COMMENT “传感器ID“, dt DATE NOT NULL COMMENT “数据日期用于分区“, timestamp DATETIME NOT NULL COMMENT “精确时间戳“, temperature DECIMAL(5,2) COMMENT “水温(摄氏度)“, salinity DECIMAL(5,3) COMMENT “盐度(PSU)“, ph DECIMAL(3,2) COMMENT “酸碱度“, location POINT COMMENT “地理位置“ ) DUPLICATE KEY(sensor_id, dt, timestamp) -- 指定排序列 PARTITION BY RANGE(dt) ( PARTITION p202405 VALUES LESS THAN (“2024-06-01“), PARTITION p202406 VALUES LESS THAN (“2024-07-01“), PARTITION p202407 VALUES LESS THAN (“2024-08-01“) ) DISTRIBUTED BY HASH(sensor_id) BUCKETS 8 PROPERTIES ( “replication_num“ “1“ -- 副本数单机设置为1 ); -- 创建一张聚合表Aggregate Key 模型用于快速查询每个传感器的最新读数。 CREATE TABLE IF NOT EXISTS sensor_data_latest ( sensor_id INT NOT NULL COMMENT “传感器ID“, dt DATE NOT NULL COMMENT “数据日期“, latest_timestamp DATETIME MAX COMMENT “最新数据时间“, latest_temperature DECIMAL(5,2) REPLACE COMMENT “最新水温“, latest_salinity DECIMAL(5,3) REPLACE COMMENT “最新盐度“, latest_ph DECIMAL(3,2) REPLACE COMMENT “最新酸碱度“ ) AGGREGATE KEY(sensor_id, dt) DISTRIBUTED BY HASH(sensor_id) BUCKETS 4 PROPERTIES ( “replication_num“ “1“ );关键点解释DUPLICATE KEY仅影响数据在底层存储的排序方式用于优化范围查询。AGGREGATE KEYREPLACE相同 Key 的数据导入时新数据会替换旧数据从而始终保持每个传感器的最新状态。PARTITION BY RANGE未来可以方便地添加新分区ALTER TABLE ... ADD PARTITION或删除旧分区。4.3 多种方式导入“海”量数据Doris 支持丰富的数据导入方式这里介绍最常用的三种。方式一Broker Load适用于 HDFS 或云存储上的大规模批量数据假设我们将原始 CSV 数据文件上传到了 HDFS 的/user/ocean/data/路径下。LOAD LABEL ocean_monitor.label_20240527_01 ( DATA INFILE(“hdfs://your-namenode:8020/user/ocean/data/*.csv“) INTO TABLE sensor_data_detail COLUMNS TERMINATED BY “,“ FORMAT AS “csv“ (sensor_id, timestamp, temperature, salinity, ph, location) SET ( dt DATE(timestamp) -- 从timestamp列推导出分区列dt ) ) WITH BROKER “your_broker_name“ PROPERTIES ( “timeout“ “3600“ );通过SHOW LOAD WHERE LABEL ‘label_20240527_01’;查看导入状态。方式二Stream Load适用于本地文件或程序流式写入HTTP协议使用curl命令或程序 SDK 直接推送数据这是实时导入的常用方式。curl -u root:your_password -H “format: csv“ -H “column_separator:,“ -T /path/to/local/data.csv http://fe_host:8030/api/ocean_monitor/sensor_data_detail/_stream_load可以在PROPERTIES中指定“exec_mem_limit““2147483648“等参数控制导入内存。方式三Routine Load持续消费 Kafka 等消息队列中的数据这是实现流式数据实时入库的核心功能。CREATE ROUTINE LOAD ocean_monitor.kafka_ocean_load ON sensor_data_detail COLUMNS(sensor_id, timestamp, temperature, salinity, ph, location, dtDATE(timestamp)), COLUMNS TERMINATED BY “,“ PROPERTIES ( “desired_concurrent_number““3“, “max_batch_interval““20“, “max_batch_rows““200000“, “max_batch_size““104857600“ ) FROM KAFKA ( “kafka_broker_list“ “kafka_host1:9092,kafka_host2:9092“, “kafka_topic“ “ocean_sensor_topic“, “property.group.id“ “doris_consumer_group“ );4.4 执行分析查询与验证数据导入后即可体验 Doris 的快速分析能力。查询1查询特定传感器在某个时间段内的详细数据。SELECT * FROM sensor_data_detail WHERE sensor_id 1001 AND dt ‘2024-05-01‘ AND timestamp BETWEEN ‘2024-05-01 10:00:00‘ AND ‘2024-05-01 12:00:00‘ ORDER BY timestamp DESC LIMIT 100;Doris 会利用分区和排序列进行快速过滤和排序。查询2聚合分析计算每个传感器当天的平均水温和最大盐度。SELECT sensor_id, dt, AVG(temperature) AS avg_temp, MAX(salinity) AS max_salinity, COUNT(*) AS data_count FROM sensor_data_detail WHERE dt ‘2024-05-27‘ GROUP BY sensor_id, dt ORDER BY avg_temp DESC;查询3从聚合表瞬时获取所有传感器的最新状态。SELECT * FROM sensor_data_latest WHERE dt CURRENT_DATE();由于使用了 REPLACE 聚合这个查询会非常快适合做监控大盘。查询4空间数据查询。查询某个海域范围内的所有传感器。SELECT sensor_id, ST_AsText(location) as point, temperature FROM sensor_data_detail WHERE dt ‘2024-05-27‘ AND ST_Contains(ST_PolygonFromText(‘POLYGON((120 30, 121 30, 121 31, 120 31, 120 30))‘), location);Doris 支持丰富的 GIS 函数非常适合处理带地理位置信息的数据。4.5 通过物化视图进行查询加速对于非常复杂或高频的聚合查询可以创建物化视图进行预计算。-- 创建一个存储每小时、每个传感器平均指标的物化视图 CREATE MATERIALIZED VIEW sensor_hourly_mv AS SELECT sensor_id, DATE_TRUNC(‘hour‘, timestamp) as hour_time, AVG(temperature) as avg_temp, AVG(salinity) as avg_salinity, COUNT(*) as cnt FROM sensor_data_detail GROUP BY sensor_id, DATE_TRUNC(‘hour‘, timestamp); -- 查询时优化器会自动匹配并路由到物化视图 SELECT sensor_id, hour_time, avg_temp FROM sensor_data_detail -- 注意查询的还是原表 WHERE hour_time ‘2024-05-27 10:00:00‘ ORDER BY avg_temp DESC;5. 常见问题与排查思路在 Doris 使用过程中可能会遇到以下典型问题。问题现象可能原因排查思路与解决方案导入失败报错Tablet writer write failedBE 节点磁盘空间不足单个 Tablet 数据量过大副本数设置不合理。1. 检查 BE 的storage_root_path磁盘使用率 (df -h)。2. 通过SHOW TABLET FROM table_name查看 Tablet 状态和大小。3. 考虑增加 Bucket 数量或优化分区策略使数据分布更均匀。查询速度慢EXPLAIN显示未进行分区裁剪查询条件中的分区列使用了函数或不符合分区格式。1. 确保WHERE条件直接使用分区列如dt而不是对其施加函数如DATE(timestamp)。2. 在表设计时尽量让常用查询条件直接对应分区列。Routine Load消费 Kafka 数据积压导入速度跟不上 Kafka 生产速度max_batch_*参数设置过小。1. 增加desired_concurrent_number提高并发任务数。2. 适当调大max_batch_interval,max_batch_rows,max_batch_size。3. 检查 BE 节点负载和网络带宽。内存超限错误Memory exceed limit复杂查询或导入任务申请内存超过限制。1. 对于查询可通过SET exec_mem_limitxxx;会话级调整或优化 SQL如减少全表扫描。2. 对于导入在LOAD语句的PROPERTIES中设置“exec_mem_limit““2147483648“2GB。3. 检查 BE 配置mem_limit和storage_page_cache_limit。SHOW BACKENDS显示 BE 状态异常BE 进程挂掉网络不通心跳失败。1. 登录 BE 服务器检查进程 ps aux6. 最佳实践与工程建议表设计是性能的基石前缀索引DUPLICATE/UNIQUE KEY列的顺序至关重要。将查询中最常用来过滤和排序的列放在前面以充分利用前缀索引。分区与分桶时序数据必须分区。分桶列选择高基数列如ID避免数据倾斜。分桶数 BE节点数 * 磁盘数 * 2或3。数据类型使用最精确、最小的数据类型。例如能用INT就不用BIGINT能用VARCHAR(20)就不用STRING。数据导入策略小批量高频 vs 大批量低频Stream Load适合秒/分钟级延迟Broker Load适合小时/天级T1数据。根据业务对实时性的要求混合使用。导入原子性一个导入任务Label内的数据要么全部成功要么全部失败。利用这个特性保证数据一致性。监控导入任务定期检查information_schema.loads表监控导入成功率、耗时和流量。查询优化善用EXPLAIN在复杂查询前使用EXPLAIN或EXPLAIN GRAPH查看执行计划关注SCAN行数、是否命中分区/索引、聚合节点开销。**避免 SELECT ***明确列出所需列减少网络传输和内存占用。物化视图的权衡物化视图用空间换时间。只为最关键、最耗时的聚合查询创建并注意其维护成本。集群管理与运维容量规划提前规划存储和计算资源。存储量 ≈ 原始数据量 * 副本数 * 压缩比约0.3-0.5。内存要充足用于查询和导入。监控告警集成 Prometheus Grafana监控集群健康度FE/BE状态、查询延迟、导入吞吐、磁盘/内存使用率等核心指标。备份与恢复定期使用BACKUP命令将数据快照到对象存储如S3、OSS并测试RESTORE流程。对于分区表可以结合ALTER TABLE DROP PARTITION进行历史数据清理。安全与权限生产环境务必修改默认的 root 密码。使用CREATE USER和GRANT按需分配数据库、表的读写权限遵循最小权限原则。如果通过公网访问 FE考虑设置防火墙或使用代理。从模拟的海洋传感器数据接入到完成 Doris 集群部署、表结构设计、多种方式的数据导入再到执行复杂的时空聚合查询和性能调优我们走完了一个典型的 OLAP 系统构建流程。Doris 凭借其极致的性能、简洁的架构和 MySQL 协议兼容性确实为处理“海”量数据分析提供了优秀的解决方案。在实际项目中建议先从核心业务场景的一两张表开始试点逐步积累运维经验。接下来可以进一步探索 Doris 的向量化计算、外部表如查询 Hive 数据、以及更复杂的多表 Join 优化等高级特性让数据真正成为驱动业务决策的“海洋”。
RELATED READING

延伸阅读

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