Doris 导入 1 TB 数据耗时 4 小时,如何通过并行度调优降到 1 小时

解读

  1. 业务场景:离线或准实时批量导入,数据量 1 TB,当前端到端耗时 4 小时,目标 1 小时,提速 4 倍。
  2. 性能瓶颈定位思路:先确认“慢”发生在哪一段——客户端拆数据、网络传输、FE 规划、BE 写入、Compaction、副本同步、还是导入作业串行排队。
  3. 并行度调优核心:把“串行”改成“并行”,把“小并发”改成“大并发”,同时避免磁盘 I/O、CPU、内存、网络任何一项先触顶,导致线程空转或反压。
  4. 国内面试常考“量化”:给出可落地的参数、公式、验证指标,而不是泛泛而谈“调大即可”。

知识点

  1. Doris 导入模型
    • Stream Load:同步,小文件高并发,默认 1 个作业单线程。
    • Broker Load:异步,走 HDFS/S3,多 Broker 进程并行扫描。
    • Spark Load:外部 Spark 作业预计算分桶文件,再批量 move 到 BE。
    • Routine Load:Kafka 流式,分区数决定并行上限。
  2. 并行度层级
    • 作业级:同时提交的导入作业数 ≤ max_running_txn_num_per_db。
    • Fragment 级:一个作业拆成几个 Fragment,由 FE 的 parallel_fragment_exec_instance_num 决定。
    • Scan 级:Broker Load 的 broker_scan_node_pool_size、单个文件 split size。
    • 写入级:BE 的 tablet_writer_threads、flush_threads、delta_writer_threads。
    • Compaction 级: cumulative_compaction_threads、base_compaction_threads。
  3. 资源上限
    • 单 BE 磁盘带宽 ≈ 200 MB/s(SATA)/ 550 MB/s(SSD)/ 1 GB/s(NVMe)。
    • 万兆网卡 1.25 GB/s,千兆 125 MB/s。
    • CPU 核数决定可同时跑的 fragment 实例数。
  4. 量化公式
    目标吞吐 = 1 TB × 1024 GB / 1 h = 284 MB/s。
    单并发吞吐 20 MB/s 时,理论最小并发 = 284 / 20 ≈ 14;留 30 % 冗余,取 20 并发。
  5. 监控指标
    • BE 面板:BytesWrittenPerSecond、LoadAvg、DiskUtil%、NetworkIn。
    • FE 面板:RunningLoadJobs、TxnStatus、QueryLatency。
    • 自定义埋点:作业提交时间 T0、数据落盘时间 T1、副本可见时间 T2。

答案

第一步:确认瓶颈

  1. 查看 FE audit 日志,若 80 % 时间卡在“LOADING”状态,说明 BE 写入慢;若卡在“QUEUEING”状态,说明作业排队。
  2. 对比 BE 监控,若 DiskUtil 持续 ≥ 85 %,则磁盘 I/O 是瓶颈;若 NetworkIn 打满网卡,则网络瓶颈;若 CPU idle 仍 > 30 %,说明线程数不足。

第二步:选导入方式
1 TB 属于大文件,优先 Broker Load 或 Spark Load,可充分利用多 Broker/Spark Executor 并行扫描,避免单 Stream Load 的 1 线程限制。

第三步:调作业级并行

  1. 调大 FE 配置 max_running_txn_num_per_db = 50(默认 100,可再上调,但需保证内存足够)。
  2. 把 1 TB 文件按 256 MB 切分,得到约 4000 个分片;每 200 个分片作为一个作业,共 20 个作业同时提交,避免单作业元数据过大。

第四步:调 Fragment 级并行

  1. 设置 parallel_fragment_exec_instance_num = BE 节点数 × 每节点 CPU 核数 × 0.8(留 20 % 给 Compaction 和查询)。
    例:10 台 BE,每台 16 核 → 128 × 0.8 ≈ 100 实例。
  2. Broker Load 增加 broker_scan_node_pool_size = 100,与上面实例数对齐。

第五步:调写入线程

  1. BE 配置文件 be.conf:
    tablet_writer_threads = 16
    flush_threads = 8
    delta_writer_threads = 8
  2. 若磁盘为 SSD,可再上调到 24/12/12,并开启 enable_vertical_compaction = true,降低写放大。

第六步:抑制 Compaction 抢占

  1. 导入期间临时调大 cumulative_compaction_check_interval_seconds = 120(默认 10 s),降低后台任务频率。
  2. 设置 disable_compaction_until_seconds = 1800,让数据先快速导入,之后再集中做 Compaction。

第七步:网络优化

  1. 客户端与 Broker 均使用万兆网卡,关闭交换机的 flow-control,开启 jumbo frame 9000。
  2. 若跨机房,先通过 DistCp 把数据搬到 Doris 集群同交换机下的 HDFS,再执行 Broker Load,节省 50 % 带宽。

第八步:验证

  1. 20 并发作业同时跑,单作业吞吐 15 MB/s,总吞吐 300 MB/s,1 TB 理论耗时 3417 s ≈ 57 min,满足 < 1 h。
  2. 监控 BE DiskUtil 稳定在 70 % 左右,CPU 60 %,网络 800 Mb/s,无反压,证明并行度与资源匹配。

拓展思考

  1. 如果数据量再扩大到 10 TB,单集群磁盘带宽会成为硬瓶颈,此时需:
    • 采用 Spark Load 预分桶,直接生成 Tablet 文件,跳过 Stream 写入路径,可把极限吞吐提升到 1 GB/s 级别;
    • 或者采用 Doris Multi-Cluster 动态扩副本,先把数据导入到临时 3 倍扩容集群,再缩容回正常副本,利用云硬盘弹性 IOPS。
  2. 实时链路场景下,Routine Load 的并行度由 Kafka Topic 分区数决定,分区不足时即使调大 max_routine_load_task_concurrent_num 也无法提速,需要提前按“吞吐 × 时间 / 单分区吞吐”反推分区数,并在 Kafka 端做 re-partition。
  3. 并行度并非越高越好,超过磁盘 IOPS 上限后,随机写会退化成顺序写抖动,出现“线程越多越慢”的反效果,因此每次调优后必须用控制变量法:固定并发梯度 5→10→20→40,记录吞吐曲线,找到拐点。
  4. 国内金融、运营商生产环境常开审计与镜像备份,导致磁盘实际可用带宽再打 8 折,性能测试报告必须注明“生产预留系数”,避免上线后 SLA 失守。