原创第 10 篇

更多 Flink / Spark / Doris / Clickhouse / Paimon 实战文章,请关注公众号「junzibuqi_x」

摘要

  • 背景:src  层“天-小时”二级分区表的归档任务,原方案完全串行,预计执行约 15 天。

  • 思路:采用“表间串行、表内天级并行、天内小时串行”的三级执行架构;并通过基于数据量的动态并发、Spark UI 可观测性、空分区跳过与资源配置优化,改造为可并发、可监控、可降级的技术方案。

  • 产出:归档任务从笨重的长跑转为稳健的并发工程。文内给出关键代码片段与优化方向,便于复制与落地。

环境与前提

  • 技术栈:PySpark 3.3.1

  • 任务对象:src  层二级分区(dt × hour)。

  • 设计原则:稳定优先(表间串行)、资源平衡(天级并行)、可观测(JobGroup/Description)、可降级(慢表串行)。

  • 风险约束:对“慢表”强制串行,避免并发导致驱动竞争与资源抢占。

架构总览

  • 表间串行:按顺序处理每张表,避免跨表并发互相抢占资源。

  • 表内天级并行:同一表在“天维度”并行(线程池),兼顾速度与稳定。

  • 天内小时串行:同一天的 24 小时按序处理,降低任务粒度并发导致的资源竞争。

关键设计与核心代码

1) 基于数据量的动态并发

通过估算“某表在日期范围内的最大天数据量”,动态决定天级并发数,同时保留手动指定关键大表的“慢表串行”的兜底策略。

def_decide_workers_by_size(self, size_gb: float) -> int:
    """并发数规则:<100G -> 10;100~500G -> 5;>500G -> 1"""
    if size_gb < 100:
        return10
    elif size_gb <= 500:
        return5
    else:
        return1

def_decide_workers_for_table(self, table: str, dates: List[str]) -> int:
    """计算该表在给定日期范围内的最大天数据量,并据此决定并发数"""
    max_size = 0.0
    for dt in dates:
        s = self._estimate_day_size_gb(table, dt)
        max_size = max(max_size, s)
    workers = self._decide_workers_by_size(max_size)
    if table in self.fifo_tables:  # 慢表强制串行
        workers = 1
    logging.info(f"[CONCURRENCY] table={table} max_day_size≈{max_size:.2f} GB -> workers={workers}")
    return workers

2) HDFS 天级数据量估算

通过  FileSystem.getContentSummary  获取天分区目录体量(GB),为并发策略提供依据。

def_estimate_day_size_gb(self, table: str, dt: str) -> float:
    """使用 HDFS FileSystem 估算某天分区目录的总大小(GB)"""
    jvm = self.spark._jvm
    jsc = self.spark.sparkContext._jsc
    fs = jvm.org.apache.hadoop.fs.FileSystem.get(jsc.hadoopConfiguration())
    path = jvm.org.apache.hadoop.fs.Path(self._get_day_dir(table, dt))
    try:
        if fs.exists(path):
            summary = fs.getContentSummary(path)
            size_bytes = summary.getLength()
            return float(size_bytes) / (1024.0 ** 3)
        else:
            return0.0
    except Exception as e:
        logging.warning(f"Estimate size failed: table={table} dt={dt}: {e}")
        return0.0

3) Spark UI 可观测性增强

为线程与作业设置明确的  JobGroup  和  JobDescription,便于观测与定位。

def_thread_initializer(self, table: str):
    # 线程初始化:设置可继承的 Job 描述
    base_desc = f"ArchiveRunner for table={table}"
    sc = self.spark.sparkContext
    sc.setLocalProperty("spark.job.description", base_desc)

def_archive_hour(self, table: str, dt: str, hour: str):
    """归档单小时分区:INSERT OVERWRITE DIRECTORY + Gzip 压缩"""
    target_dir = f"{self.hdfs_base_dir}/{table}/dt={dt}/hour={hour}"
    sc = self.spark.sparkContext
    job_group = f"archiver:{table}:{dt}"
    job_desc = f"[{table}] dt={dt} hour={hour}"
    sc.setLocalProperty("spark.job.description", job_desc)
    sc.setJobGroup(job_group, job_desc, True)

    sql = f"""
    INSERT OVERWRITE DIRECTORY '{target_dir}'
    SELECT /*+ REPARTITION(1) */ * FROM {self.database}.{table}
    WHERE dt='{dt}' AND hour='{hour}'
    """
    logging.info(f"[ARCHIVE] table={table} dt={dt} hour={hour} -> {target_dir}")
    self.spark.sql(sql)

4) 空分区跳过机制

避免对无数据的小时分区做无效重写,节省资源与时间。

def_partition_has_data(self, table: str, dt: str, hour: str) -> bool:
    """判断该小时分区是否有数据,避免无意义的重写"""
    sql = f"SELECT COUNT(1) AS c FROM {self.database}.{table} WHERE dt='{dt}' AND hour='{hour}'"
    try:
        cnt = self.spark.sql(sql).collect()[0]["c"]
        return cnt and cnt > 0
    except Exception as e:
        logging.warning(f"Count check failed for {table} dt={dt} hour={hour}: {e}")
        returnFalse

5) 会话配置与资源优化

设置并行度、输出压缩与动态分区覆盖,确保稳定性与资源均衡。确保调度器使用公平调度(FAIR),因为默认的 FIFO 会导致并发策略无效,导致长作业阻塞短作业。

def_configure_session(self):
    """Spark 会话配置优化"""
    // ... existing code ...
    // 启用公平调度(FAIR),避免长作业阻塞短作业
    self.spark.conf.set("spark.scheduler.mode", "FAIR")

    # 输出压缩设置
    self.spark.sql("SET hive.exec.compress.output=true")
    self.spark.sql(f"SET mapreduce.output.fileoutputformat.compress.codec={self.compression_codec}")
    self.spark.sql("SET mapreduce.output.fileoutputformat.compress.type=BLOCK")

    # 动态分区覆盖
    self.spark.sql("SET spark.sql.sources.partitionOverwriteMode=dynamic")

6) 三级执行策略的具体落地

  • “小时串行”的日内执行,避免粒度并发导致 Driver 竞争

  • “天级并行”的线程池控制,配合并发上限与失败隔离

def_archive_one_day(self, table: str, dt: str):
    """串行处理一天的24小时,避免过多并发作业导致 Driver 竞争"""
    logging.info(f"[DAY-START] table={table} dt={dt}")
    ok, skip = 0, 0
    for h in self.hours:
        ifnot self._partition_has_data(table, dt, h):
            logging.info(f"[SKIP] table={table} dt={dt} hour={h} (no data)")
            skip += 1
            continue
        try:
            self._archive_hour(table, dt, h)
            ok += 1
        except Exception as e:
            logging.error(f"[FAIL] table={table} dt={dt} hour={h}: {e}")
    logging.info(f"[DAY-DONE] table={table} dt={dt} archived={ok} skipped={skip}")

defarchive_table(self, table: str, dates: List[str]):
    """针对某张表:按“天级并发(线程池)+ 小时串行”的策略归档"""
    workers = self._decide_workers_for_table(table, dates)
    logging.info(f"[TABLE-START] {table} dates={len(dates)} days, max_concurrent_days={workers}")
    with ThreadPoolExecutor(
        max_workers=workers,
        initializer=self._thread_initializer,
        initargs=(table,)
    ) as executor:
        futures = {executor.submit(self._archive_one_day, table, dt): dt for dt in dates}
        for f in as_completed(futures):
            dt = futures[f]
            try:
                f.result()
            except Exception as e:
                logging.error(f"[DAY-ERROR] table={table} dt={dt}: {e}")
    logging.info(f"[TABLE-DONE] {table}")

defrun(self, tables: List[str], dates: List[str]):
    """按顺序处理每张表,避免跨表并发对资源的争抢"""
    for table in tables:
        self.archive_table(table, dates)

7) 命令行接口与使用示例

支持  YYYYMM  月份或日期范围,小时可选,默认全量。

defparse_args():
    parser = argparse.ArgumentParser(description="Spark SQL 每小时分区归档(按天并行、按小时串行)")
    group = parser.add_mutually_exclusive_group(required=True)
    group.add_argument("--month", type=str, help="月份(YYYYMM),例如 202507")
    group.add_argument("--start", type=str, help="开始日期(YYYYMMDD)")
    parser.add_argument("--end", type=str, help="结束日期(YYYYMMDD),与 --start 搭配使用")
    parser.add_argument("--tables", type=str, default=",".join(_default_tables()),
                        help="逗号分隔的表名列表,默认为预置清单")
    parser.add_argument("--hdfs-base-dir", type=str,
                        default="hdfs://hdfs-cluster/hive/warehouse/jaco/jaco_src",
                        help="HDFS 基路径,例如 hdfs://hdfs-cluster/hive/warehouse/jaco/jaco_src")
    parser.add_argument("--database", type=str, default="jaco_src", help="Hive 库名,默认 jaco_src")
    parser.add_argument("--max-concurrent", type=int, default=4, help="每张表的天级并发数")
    parser.add_argument("--shuffle-partitions", type=int, default=200, help="Spark SQL 并行度")
    parser.add_argument("--hours", type=str, default="", help="逗号分隔的小时(00..23),不传则默认全量")
    parser.add_argument("--compression-codec", type=str,
                        default="org.apache.hadoop.io.compress.GzipCodec", help="输出压缩 codec,默认 Gzip")
    return parser.parse_args()

defmain():
    args = parse_args()
    if args.month:
        dates = SparkHourlyPartitionArchiver._parse_month_dates(args.month)
    else:
        ifnot args.start ornot args.end:
            raise ValueError("使用日期范围时,必须同时提供 --start 和 --end")
        dates = SparkHourlyPartitionArchiver._parse_range_dates(args.start, args.end)

    hours = None
    if args.hours:
        hours = [h.strip() for h in args.hours.split(",") if h.strip()]
    tables = [t.strip() for t in args.tables.split(",") if t.strip()]

    archiver = SparkHourlyPartitionArchiver(
        hdfs_base_dir=args.hdfs_base_dir,
        database=args.database,
        hours=hours,
        max_concurrent_days=args.max_concurrent,
        shuffle_partitions=args.shuffle_partitions,
        compression_codec=args.compression_codec,
    )

    logging.info(f"Tables={tables}")
    logging.info(f"Dates={dates[0]} .. {dates[-1]} ({len(dates)} days)")
    logging.info(f"Hours={archiver.hours}")
    archiver.run(tables, dates)

if __name__ == "__main__":
    main()

效果复盘(建议记录与展示)

  • 执行时长:从“15 天串行”改造为“可并发日级执行”,缩短幅度取决于数据分布与集群资源。

  • 并发峰值与资源占用:记录线程池并发数、Driver 内存峰值、网络带宽使用情况。

  • 失败与重试:统计小时分区失败率与重试成功率,验证稳定性与隔离性。

  • 空分区跳过:跳过比例与节省资源评估。

  • 压缩比与存储节省:Gzip 压缩的平均压缩比与空间节省。

方案亮点

  • 动态并发 + 慢表串行:兼顾效率与稳健,避免“一刀切”的资源过载。

  • 可观测性优先:JobGroup/Description 全链路标识,Spark UI 可视化明确。

  • 空分区跳过:减少无效重写,节省 IO 与计算。

  • 资源配置得当:并行度、压缩与覆盖策略一致,降低失败面。

  • 工程化心智:能跑、能看、能调、能迭代,具备推广价值。

结语

将“长跑型归档任务”并行化,是数仓运维走向稳健与高效的必经之路。通过动态并发、可观测性与降级策略的组合,这套方案既能提速,也能控风险。

从 2 小时到 3 分钟:Spark SQL 多维分析性能优化实战

一夜堆积 4 万条任务, ClickHouse 和我都快顶不住了

告别直播卡顿:我们如何用 Flink SQL 打造秒级监控大盘

别再无脑用 COUNT DISTINCT 了!90%的 Spark 工程师都踩过这个坑

ClickHouse 宕机排查实录:一条异常 SQL 如何让内存飙到爆

基于 Spark Thrift Server 构建统一 SQL 查询服务实践

Spark AQE 优化篇 ①:小文件治理

更多推荐