ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Doris与Hadoop生态融合实践与优化策略

2026/9/16 12:54:05 拓冰建站 浏览量
Doris与Hadoop生态融合实践与优化策略 1. 为什么需要将Doris融入Hadoop生态第一次接触Doris时我就被它的实时分析能力惊艳到了。这个由百度开源的MPP数据库能在秒级响应复杂的OLAP查询与我们团队长期使用的Hadoop生态形成了鲜明对比。但很快我就发现单纯用Doris替代Hive或Impala并不现实——现有PB级的历史数据都安静地躺在HDFS上ETL流程也早已围绕Hadoop构建完整。于是如何让Doris与Hadoop生态和谐共处就成了我们必须解决的现实问题。经过半年多的实践验证我们摸索出了一套行之有效的融合方案。Doris作为实时分析引擎处理热数据Hadoop继续承担批量计算和冷数据存储两者通过精心设计的数据通道实现无缝衔接。这种架构不仅保留了Hadoop处理海量数据的可靠性还获得了Doris带来的亚秒级查询体验。最让我意外的是某些场景下两者的协同效应甚至产生了112的效果。2. 核心组件对接方案解析2.1 与HDFS的深度集成Doris原生支持通过Broker Load方式读取HDFS数据但默认配置在TB级数据传输时表现不佳。我们通过以下优化实现了稳定高效的数据同步-- 优化后的Broker Load示例 LOAD LABEL db1.label1 ( DATA INFILE(hdfs://namenode:8020/path/to/file/*) INTO TABLE target_table FORMAT AS parquet ) WITH BROKER hdfs_broker ( usernamehadoop, passwordyour_password, dfs.nameservicesyour_nameservice, dfs.ha.namenodes.your_nameservicenn1,nn2, dfs.namenode.rpc-address.your_nameservice.nn1namenode1:8020, dfs.namenode.rpc-address.your_nameservice.nn2namenode2:8020, dfs.client.failover.proxy.provider.your_nameserviceorg.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider ) PROPERTIES ( timeout 3600, max_filter_ratio 0.1, exec_mem_limit 21474836480, strict_mode true );关键配置经验对于HA集群必须配置nameservice相关参数否则单节点故障会导致作业失败exec_mem_limit需要根据数据量调整默认4GB在处理大文件时经常OOM建议开启strict_mode避免脏数据导致整个作业失败踩坑记录曾因未设置timeout导致一个3TB的加载任务在网络波动时无限重试最终阻塞了整个集群的调度队列。现在我们会根据数据量合理设置超时阈值通常按1小时/TB估算。2.2 与Hive Metastore的元数据同步为了让Doris能直接查询Hive表我们采用了External Table方式对接。但在生产环境发现了几个关键问题分区表的新增分区需要手动refresh表结构变更不会自动同步缺乏统一的权限管控解决方案是开发元数据同步服务核心逻辑如下// 伪代码展示同步流程 public void syncHiveToDoris(String dbName) { // 获取Hive所有表 ListHiveTable hiveTables hiveClient.getAllTables(dbName); // 遍历处理每张表 for (HiveTable hiveTable : hiveTables) { // 检查Doris中是否存在对应表 DorisTable dorisTable dorisClient.getTable(dbName, hiveTable.getName()); if (dorisTable null) { // 新建外部表 dorisClient.createExternalTable(convertSchema(hiveTable)); } else { // 对比schema差异 SchemaDiff diff compareSchema(hiveTable, dorisTable); if (diff.hasChange()) { // 执行ALTER TABLE dorisClient.alterTable(dbName, hiveTable.getName(), diff.getDdl()); } } // 处理分区差异 if (hiveTable.isPartitioned()) { syncPartitions(hiveTable, dorisTable); } } }实际运行中我们设置了每小时全量同步事件触发增量同步的混合机制将元数据延迟控制在5分钟以内。3. 实时数据流架构设计3.1 Kafka作为数据枢纽的实践我们设计的实时管道架构如下[业务数据库] - (CDC) - [Kafka] - [Flink] - [Doris] │ └─ [Spark] - [HDFS] - [Hive]Flink作业的双写配置示例// Flink双写配置示例 final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 消费Kafka数据 KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka:9092) .setTopics(user_events) .setDeserializer(new SimpleStringSchema()) .build(); // Doris Sink配置 DorisSink.BuilderString dorisSinkBuilder DorisSink.builder() .setFenodes(doris-fe:8030) .setUsername(admin) .setPassword() .setTableIdentifier(db1.table1) .setSerializer(new JsonDebeziumSerializer()); // HDFS Sink配置 StreamingFileSinkString hdfsSink StreamingFileSink .forRowFormat(new Path(hdfs://path/to/raw), new SimpleStringEncoder()) .withBucketAssigner(new EventTimeBucketAssigner()) .build(); // 构建拓扑 env.fromSource(source, WatermarkStrategy.noWatermarks(), Kafka Source) .process(new EventParser()) .addSink(dorisSinkBuilder.build()) .name(Doris Sink); env.fromSource(source, WatermarkStrategy.noWatermarks(), Kafka Source) .process(new EventParser()) .addSink(hdfsSink) .name(HDFS Sink);关键调优参数Flink checkpoint间隔设置为1分钟兼顾可靠性和性能Doris sink的buffer.size100MB, buffer.count3HDFS滚动策略按1小时或1GB触发3.2 批流一体数据一致性保障我们遇到过最棘手的问题是同样的查询条件实时看板和历史报表显示的结果不一致。根本原因在于流处理中的迟到数据批处理作业的调度延迟两端计算逻辑的细微差异解决方案是引入一致性校验机制-- Doris每日与Hive数据对比的校验SQL WITH doris_stats AS ( SELECT dt, COUNT(*) AS doris_cnt, SUM(amount) AS doris_sum FROM realtime_table WHERE dt 2023-07-15 GROUP BY dt ), hive_stats AS ( SELECT dt, COUNT(*) AS hive_cnt, SUM(amount) AS hive_sum FROM offline_table WHERE dt 2023-07-15 GROUP BY dt ) SELECT d.dt, d.doris_cnt, h.hive_cnt, d.doris_sum, h.hive_sum, ABS(d.doris_cnt - h.hive_cnt) AS cnt_diff, ABS(d.doris_sum - h.hive_sum) AS sum_diff FROM doris_stats d JOIN hive_stats h ON d.dt h.dt;当差异超过阈值如0.1%时自动触发告警并由数据工程师介入排查。这套机制将数据不一致问题从原来的每周数起降低到每月不足一次。4. 混合查询优化策略4.1 冷热数据自动分层我们根据访问频率将数据划分为三个层级热数据近7天存储在Doris本地磁盘温数据7-30天存储在Doris S3冷数据30天仅存储在HDFS通过以下DDL实现自动分层CREATE TABLE user_events ( event_time DATETIME, user_id BIGINT, event_type VARCHAR(32), -- 其他字段... ) PARTITION BY RANGE(event_time) ( PARTITION p202307 VALUES LESS THAN (2023-08-01), PARTITION p202308 VALUES LESS THAN (2023-09-01) ) DISTRIBUTED BY HASH(user_id) BUCKETS 32 PROPERTIES ( storage_medium SSD, storage_cooldown_time 7 days, replication_num 3, dynamic_partition.enable true, dynamic_partition.time_unit MONTH, dynamic_partition.start -12, dynamic_partition.end 3, dynamic_partition.prefix p, dynamic_partition.buckets 32 );实际效果热数据查询P99延迟从12s降至0.8s存储成本降低60%冷数据使用HDD存储维护工作量减少自动过期旧分区4.2 联邦查询实践对于需要同时访问Doris和Hive的查询我们采用两种方案方案一Doris外表直接查询-- 创建Hive外表 CREATE EXTERNAL TABLE hive_user_profile ( user_id BIGINT, gender VARCHAR(10), age INT ) ENGINEHIVE PROPERTIES ( hive.metastore.uris thrift://hive-metastore:9083, database profile_db, table user_profile ); -- 执行联邦查询 SELECT e.user_id, p.gender, COUNT(DISTINCT e.event_type) AS event_types FROM doris_events e JOIN hive_user_profile p ON e.user_id p.user_id WHERE e.event_time BETWEEN 2023-07-01 AND 2023-07-31 GROUP BY e.user_id, p.gender;方案二通过Presto桥接-- Presto配置连接器 CREATE CATALOG doris WITH ( connector.namedoris, connection-urljdbc:mysql://doris-fe:9030, connection-useradmin, connection-password ); CREATE CATALOG hive WITH ( connector.namehive, hive.metastore.urithrift://hive-metastore:9083 ); -- 执行跨源查询 SELECT d.user_id, h.gender, COUNT(*) AS event_count FROM doris.db1.events d JOIN hive.profile_db.users h ON d.user_id h.user_id WHERE d.dt 2023-07-15 GROUP BY d.user_id, h.gender;性能对比Doris外表适合点查和小规模joinPresto方案适合复杂分析但需要额外维护Presto集群5. 运维监控体系建设5.1 关键指标监控项我们通过PrometheusGrafana构建的监控看板包含以下核心指标指标类别具体指标告警阈值采集方式资源使用FE/JVM内存使用率80%持续5分钟JMX ExporterBE磁盘使用率90%Node Exporter查询性能99分位查询延迟3sDoris Metric查询错误率1%Doris Audit Log数据同步Broker Load成功率95%Doris Job LogHDFS-Doris延迟1小时自定义脚本5.2 自动化运维脚本我们开发的几个实用脚本分区健康检查脚本#!/bin/bash # 检查Doris分区分布是否均衡 BE_LIST$(curl -s http://doris-fe:8030/api/backends | jq -r .backends[].host) for be in $BE_LIST; do echo Checking $be ... ssh $be du -h /path/to/storage | grep partition_ | sort -rh | head -n 5 done数据校验脚本# 比较Doris和Hive表数据差异 import pyhive import pymysql def compare_table(doris_conn, hive_conn, table_name, dt): # 查询Doris数据 doris_cursor doris_conn.cursor() doris_cursor.execute(fSELECT COUNT(*) FROM {table_name} WHERE dt{dt}) doris_count doris_cursor.fetchone()[0] # 查询Hive数据 hive_cursor hive_conn.cursor() hive_cursor.execute(fSELECT COUNT(*) FROM {table_name} WHERE dt{dt}) hive_count hive_cursor.fetchone()[0] diff_ratio abs(doris_count - hive_count) / hive_count return { doris_count: doris_count, hive_count: hive_count, diff_ratio: diff_ratio }6. 典型问题排查实录6.1 Broker Load卡住问题现象数据加载到90%后长时间停滞BE节点CPU持续高负载排查过程首先检查FE Master日志发现持续打印tablet writer add batch wait警告通过SHOW BACKENDS发现目标BE的磁盘IOutil持续100%进一步检查发现该BE正在执行compaction操作解决方案-- 临时调整compaction参数 SET GLOBAL enable_vertical_compaction false; SET GLOBAL cumulative_compaction_min_deltas 5; SET GLOBAL base_compaction_interval_seconds_since_last_operation 3600; -- 重启卡住的load作业 CANCEL LOAD WHERE LABEL stuck_label; -- 重新提交时增加资源限制 LOAD LABEL ... PROPERTIES (exec_mem_limit8589934592);6.2 联邦查询性能低下现象跨Doris和Hive的join查询耗时从分钟级突增到小时级根因分析执行EXPLAIN发现Hive表缺少分区裁剪检查发现Hive表新增了分区但Doris外表未refresh导致Doris将整个Hive表数据拉取到本地再过滤优化方案-- 手动刷新元数据 REFRESH EXTERNAL TABLE hive_catalog.db1.table1; -- 为常用查询创建物化视图 CREATE MATERIALIZED VIEW hive_doris_bridge DISTRIBUTED BY HASH(user_id) REFRESH ASYNC AS SELECT h.user_id, h.profile_data, d.last_active_time FROM hive_db.users h JOIN doris_db.user_events d ON h.user_id d.user_id;经过这些优化同类查询性能恢复到原有水平P99延迟从43分钟降至2.7秒。