Spark 数仓分区归档优化:从 15 天到高效并发实践
原创第 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 如何让内存飙到爆
更多推荐









所有评论(0)