
凌晨三点被电话叫醒不是家里出事了是集群里跑得好好的离线任务整批挂掉数据管道断了早晨八点的报表注定开天窗。这是我刚做数据工程那几年最常见的场景也是后来下定决心把整套大数据体系往云原生架构上搬的直接原因。我讲的这个云原生架构不是简单把 Hadoop 生态照搬到容器里跑而是从资源调度、任务编排、存储计算分离、弹性扩缩容、可观测性这几个维度把数据工程的基础设施重做一遍。这篇文章把我这两年在一套实际落地的数据平台上积累的经验、踩过的坑、调过的参完整写出来。做数据开发、平台运维或者正在准备大数据相关毕业设计和竞赛比如网约车数据清洗、数据分析、可视化这类项目的同学都能从这里找到可以直接抄作业的参考。1. 内容整体设计与思路拆解1.1 从物理机到云原生的动机我们到底在解决什么问题传统大数据平台长什么样做过的人心里都有数。机房里的物理机装好 Hadoop、Spark、Hive搭一个 YARN 集群然后所有计算任务都往 YARN 上扔。这套架构本身没有什么问题数据量真的很大的时候它的稳定性是被验证过的。但它的短板恰恰出在“不太大的时候”和“变化很快的时候”。我做数据工程那几年被这些问题反复折磨资源利用率低。按业务波峰预估机器数量结果大部分时间集群是闲着的。波峰一过几千核的 CPU 就那么空转还要承担机器折旧和电费。扩容速度慢。业务线突然要加一个实时计算任务需要给 YARN 队列单独划资源申请机器的流程走 3 天都不夸张。数据不等人业务也不等人。环境和依赖地狱。A 团队的任务需要 Spark 3.2B 团队的任务需要 Spark 2.4跑在同一批机器上部署版本冲突、依赖冲突是常事。更不要说 Python 版本、Java 版本、各种底层库一言不合就跑出莫名其妙的问题。故障域太大。一个节点上跑着几十个任务一个任务把内存打爆直接连累同一节点的其他任务。排查问题的时候需要登录到一台一台的机器上看日志效率低到让人绝望。云原生架构的核心思路就是把这些痛点逐个拆解。我把目标定得很明确让环境隔离、让资源弹性、让调度自动化、让平台可观测。用一套完整的技术体系把数据采集、数据清洗、数据分析、数据可视化各个阶段统一收敛到云端基础设施之上让整个数据平台像一台可以随时分配资源的“超级计算机”而不是一台一台需要人肉管理的物理机。1.2 架构选型背后的核心逻辑为什么是云原生而不是其他方案在做架构选型的时候我考虑过几个方向把对比结果列出来大家就明白为什么最终选择走云原生这条路。方案资源弹性环境隔离部署效率运维复杂度适用场景传统物理机 YARN差物理扩缩容差依赖冲突普遍低人工操作多高需要专门运维团队超大集群、业务规模稳定虚拟机 容器编排中受限于宿主机规划较好但性能损耗大中中中小规模过渡方案容器 Kubernetes好秒级扩缩容好镜像隔离彻底高镜像即交付中依赖 k8s 运维能力业务波动大、多租户、版本迭代快一个重要的原则是能用容器解决的全都容器化。大数据组件那么多部署方式五花八门如果不统一成镜像后面的运维痛苦是持续的。比如 Spark 的 History Server、Hive Metastore、Flink JobManager这些无状态或者半无状态的组件全部容器化部署有状态的数据存储单独走云盘或对象存储。还有一个关键决策是存储和计算分离。这句话说起来容易做起来是真的要下决心。传统 HDFS 架构下计算任务去拿数据走的是内网磁盘 IO速度飞快。但是一旦做成云原生架构最流行的方案就是 S3 兼容对象存储我用的是 MinIO计算节点从对象存储拉数据网络开销变大我就遇到过因为频繁拉取小文件导致任务慢 30% 的情况。后面通过调整文件合并策略、在客户端做缓存这个问题得到了缓解但初次改造的人往往是踩过坑才知道疼。1.3 整体技术栈映射每个部分解决什么问题给大家看我最后沉淀下来的核心技术栈每一项的定位都说得清楚容器与编排Docker Kubernetes。镜像打包所有大数据组件运行环境K8s 负责资源调度、容灾恢复、弹性扩缩容。数据存储MinIOS3 协议兼容存原始数据、清洗后的数据、结果数据关系型数据库PostgreSQL存元数据、任务调度状态。计算引擎Spark离线批处理、Flink实时流处理、Hive数据仓库查询。调度系统Apache DolphinScheduler 或者 Argo Workflows负责编排数据管道 DAG。DolphinScheduler 更贴近数据工程场景我用的就是它。数据接入Flume 采集日志Kafka 承接实时消息流Canal 同步业务库增量数据。可观测性Prometheus 采集指标Grafana 做可视化大盘Loki 收集日志Alertmanager 做告警通知。这套技术栈相互之间的配合逻辑简单用一句话总结Flume/Canal/Kafka 负责把数据搬进来Spark/Flink 负责把数据算明白Hive/MinIO 负责把数据存整齐DolphinScheduler 负责让一切任务按时跑起来Prometheus/Grafana 负责让我们看得见系统状态。2. 环境准备与部署核心选型解析2.1 硬件资源规划不要把 Kubernetes 当玩具很多第一次搞 K8s 大数据的人第一反应是先搞三台虚拟机然后发现连基础组件都跑不全。我强烈建议至少按下面这个标准起步3 台 Master 节点8C16G 起步etcd 和 API Server 都是内存敏感型组件内存给 16G 一点不浪费。如果你只有 3 台节点也想同时当 Master 和 Worker那需要给 etcd 单独配置通信端口的优先级否则高负载下 etcd 的读写延迟会让整个集群状态同步出问题。5 台以上 Worker 节点16C64G 是标准配置。Spark Executor 和 Flink TaskManager 都非常吃内存单节点内存太小会导致任务频繁 OOM而且 Kubernetes 默认的 Pod 调度策略会尽量把 Pod 打散到不同节点节点太少会造成资源碎片化。操作系统选择统一使用 Ubuntu 22.04 LTS 或 Rocky Linux 9内核版本和容器运行时兼容性比较好。存储规划SSD 做本地数据盘跑中间结果和 Spark Shuffle机械盘分布式存储放冷数据日志单独挂载一块盘避免日志写满导致的系统故障。很多人忽略一个重要细节K8s 集群的时间同步比物理机集群的要求更高。大数据组件对时间漂移极其敏感Spark 和 Flink 的任务状态、数据 watermark 都依赖精确时钟。我在最初部署时没有配置 chrony 时间同步结果 Flink 做窗口聚合的时候经常出现窗口延迟几分钟的诡异现象排查到最后才发现是各节点时间差了 12 秒。所以所有节点必须统一配置 chrony 并强制校时这个看起来不起眼的步骤能省掉后续大量的排查时间。2.2 安装层组件部署 K8s 和 Helm 的实操记录Kubernetes 部署方式我用的是 kubeadm生产单集群需要快速可控地搭建二进制方式运维成本太高。书写这个过程用kubeadm init的时候有下面几个参数一定要加上不能省kubeadm init \ --apiserver-advertise-address192.168.1.10 \ --pod-network-cidr10.244.0.0/16 \ --service-cidr10.96.0.0/12 \ --kubernetes-versionv1.28.2--pod-network-cidr这个参数的作用是给整个集群的 Pod 预设一个网段Flannel 网络组件会基于这个网段来规划子网。如果忘了设置后加网络插件会遇到路由不同的问题重来一遍成本极高。初始化完成后需要把生成的配置放到~/.kube/config这一步经常有人漏掉然后kubectl get nodes报连接拒绝其实是 kubeconfig 没配好mkdir -p $HOME/.kube sudo cp -i /etc/kubernetes/admin.conf $HOME/.kube/config sudo chown $(id -u):$(id -g) $HOME/.kube/config网络插件层面我选择 Flannel 而不是 Calico主要原因是 Flannel 足够简单VXLAN 模式对大数据流量的支持够用且不依赖额外数据库。如果业务对网络策略有强诉求再考虑 Calico否则 Flannel 能让我们少踩很多配置坑。在大数据组件部署时强烈推荐使用 Helm。比如部署 Flink Operator 和 Spark Operator直接两条命令helm repo add flink-operator https://apache.github.io/flink-kubernetes-operator helm install flink-kubernetes-operator flink-operator/flink-kubernetes-operator helm repo add spark-operator https://googlecloudplatform.github.io/spark-on-k8s-operator helm install spark-operator spark-operator/spark-operator这样做的意义在于通过 Operator 方式运行 Spark/Flink 任务可以直接用 Kubernetes 原生的方式管理计算实例生命周期。任务提交后不用关心 Driver 和 Executor 具体落在哪个节点失败时 Operator 会自动做重启和资源回收。2.3 数据中间件部署Kafka、MinIO、Hive Metastore 的关键配置Kafka 在 K8s 里部署我用的是 Strimzi Kafka Operator它能把 Kafka 集群的所有核心组件Broker、ZooKeeper、TopicOperator、UserOperator全部以 CRD 方式托管。创建一个 Kafka 集群需要写一份 YAMLapiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name:>a1.sources r1 a1.sinks k1 a1.channels c1 a1.sources.r1.type taildir a1.sources.r1.positionFile /opt/flume/position/taildir_position.json a1.sources.r1.filegroups f1 a1.sources.r1.filegroups.f1 /data/logs/order/.*\.log a1.sources.r1.fileHeader true a1.channels.c1.type memory a1.channels.c1.capacity 100000 a1.channels.c1.transactionCapacity 5000 a1.sinks.k1.type org.apache.flume.sink.kafka.KafkaSink a1.sinks.k1.kafka.bootstrap.servers kafka-cluster-kafka-bootstrap:9092 a1.sinks.k1.kafka.topic order-log a1.sinks.k1.kafka.flumeBatchSize 2000 a1.sinks.k1.kafka.producer.acks 1 a1.sinks.k1.kafka.producer.linger.ms 50这里面比较关键的是linger.ms参数它表示生产者等待一批消息的时间。设置为 50 毫秒是为了在实时性和吞吐量之间做平衡——如果设置为 0每条日志立即发送吞吐量会下降设置太大实时性变差。50 毫秒配合 2000 条的flumeBatchSize单机每秒能稳定写 2 万条以上日志到 Kafka。Canal 同步数据库变更数据的部署稍微复杂一点Canal 服务端需要连接到 MySQL 的 binlog然后解析成结构化数据推送到 Kafka。需要注意 MySQL 的binlog_format必须设置为ROW模式且binlog_row_image为FULL否则增量同步过来的数据会丢字段。3.2 离线清洗链路Spark on Kubernetes 提交任务的完整示例数据进到 Kafka 之后离线链路每天凌晨从 Kafka 或对象存储读取原始数据用 Spark 做清洗最终写入 Hive 表或 MinIO Parquet 文件。我之前写过一个订单日志清洗任务把原始 JSON 日志解析成结构化数据再过滤掉无效记录完整提交给 K8s 集群运行的方式如下spark-submit \ --master k8s://https://kubernetes.default.svc:443 \ --deploy-mode cluster \ --name order-log-cleaning \ --class com.data.cleaning.OrderLogCleaner \ --conf spark.kubernetes.authenticate.driver.serviceAccountNamespark \ --conf spark.kubernetes.container.imageregistry.internal/data/spark3.4:v1.0 \ --conf spark.kubernetes.container.image.pullPolicyAlways \ --conf spark.kubernetes.driver.pod.nameorder-clean-driver \ --conf spark.executor.instances8 \ --conf spark.executor.cores2 \ --conf spark.executor.memory4g \ --conf spark.driver.memory2g \ --conf spark.sql.shuffle.partitions64 \ --conf spark.shuffle.partitions64 \ --conf spark.dynamicAllocation.enabledfalse \ --conf spark.kubernetes.executor.podTemplateFile/conf/executor-template.yaml \ local:///opt/app/order-clean-job.jar有几个参数配得有讲究我逐个说spark.executor.instances8和spark.executor.cores2每个 Executor 分配 2 核一共 8 个实例总共 16 核计算资源。为什么不把 cores 调大Spark 中单 Executor 核数越多并发 shuffle 时对内存和网络压力越大2 核是最容易做故障隔离的配置——一个 Executor OOM 了影响的只是它的两个任务不会拖累整批任务。spark.executor.memory4g加上 overhead 内存默认 10%每个 Executor Pod 实际申请的内存接近 4.5G。留给操作系统的余量足够不会因为节点内存紧张而触发驱逐。spark.dynamicAllocation.enabledfalse定死 Executor 数量。因为离线清洗任务的数据体量是基本可预估的——日志量每天变化不大动态伸缩反而会引入调度开销而且动态分配在 K8s 上需要额外的 shuffle service配置复杂度高、收益又有限。spark.sql.shuffle.partitions64控制 shuffle 阶段的并发度。数据量大约 1TB 左右每个分区处理 16GB是 Spark 的最佳实践区间不至于分区太细导致调度开销过大。还有一个特别要注意的点spark.kubernetes.executor.podTemplateFile指定的 Pod 模板文件里面配置了 Executor 容器的资源限额、环境变量、挂载卷。比如我们要给 Executor 挂载一个用于访问 MinIO 的~/.aws/credentials直接写在 Pod 模板里避免在作业代码里硬编码密钥。3.3 实时计算链路Flink SQL 消费 Kafka 到结果写入的完整作业实时场景以网约车平台的核心需求为例实时统计每分钟每个区域的订单量、活跃司机数、平均成交时长。这个需求放在 Flink SQL 里声明式地完成CREATE TABLE order_topic ( order_id STRING, driver_id STRING, city STRING, district STRING, order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic order-log, properties.bootstrap.servers kafka-cluster-kafka-bootstrap:9092, properties.group.id flink-order-stat, scan.startup.mode latest-offset, format json ); CREATE TABLE district_order_stat ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), district STRING, order_count BIGINT, active_driver_count BIGINT, avg_duration_second DOUBLE ) WITH ( connector jdbc, url jdbc:postgresql://postgres-service:5432/order_warehouse, table-name district_order_stat, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 5s ); INSERT INTO district_order_stat SELECT TUMBLE_START(order_time, INTERVAL 1 MINUTE) AS window_start, TUMBLE_END(order_time, INTERVAL 1 MINUTE) AS window_end, district, COUNT(order_id) AS order_count, COUNT(DISTINCT driver_id) AS active_driver_count, AVG(DATE_DIFF(second, order_time, CURRENT_TIMESTAMP)) AS avg_duration_second FROM order_topic GROUP BY TUMBLE(order_time, INTERVAL 1 MINUTE), district;这段 SQL 里最重要的设计是watermark 设置。网约车订单日志实际产生的时间和服务端接收日志的时间可能会因为网络原因有几秒偏差WATERMARK FOR order_time AS order_time - INTERVAL 5 SECOND就是告诉 Flink允许最多 5 秒的乱序数据超过 5 秒的迟到数据会被丢弃。scan.startup.mode latest-offset表示只从 Kafka 最新的偏移开始消费不回溯历史。对实时大屏场景这个配置合理如果要做状态恢复、容错就需要改成earliest-offset或者指定时间戳。写入端用了 JDBC Sink重点是sink.buffer-flush.max-rows1000和sink.buffer-flush.interval5s这两个参数。它们的作用是把小批量数据攒起来再批量写入 PostgreSQL避免每条结果都开一次数据库连接。我实测下来这个配置对 PostgreSQL 的写入压力能减少 80% 以上。Flink 任务提交到 K8s 上可以通过 Flink Kubernetes Operator 定义FlinkDeployment对象以声明式方式管理任务生命周期apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: order-stat-job spec: image: registry.internal/data/flink-demo:1.0 flinkVersion: v1_18 serviceAccount: flink podTemplate: spec: containers: - name: flink-main-container env: - name: KAFKA_BROKERS value: kafka-cluster-kafka-bootstrap:9092 - name: POSTGRES_HOST value: postgres-service flinkConfiguration: taskmanager.numberOfTaskSlots: 2 parallelism.default: 4 jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 4096m state.backend: rocksdb state.checkpoints.dir: s3://flink-checkpoints/order-stat-job state.savepoints.dir: s3://flink-checkpoints/savepoints/order-stat-job execution.checkpointing.interval: 60s注意到state.backend: rocksdb这个选择很关键。在 K8s 环境下Stateful 状态的存储尽量用 RocksDB 而不是堆内存——因为 Pod 内存受限堆内存一旦被状态撑爆任务直接 OOM 崩溃。RocksDB 把状态写到本地盘配合增量 checkpoint容错能力和资源占用都更理想。execution.checkpointing.interval60s表示 Flink 每 60 秒做一次分布式快照。这个频率在“故障恢复时间”和“checkpoint 存储成本”之间取了平衡。频率太高对象存储写入压力大太低重启时状态回滚太多会出现几分钟的重复计算。3.4 可视化环节Flask ECharts 展示实时指标的方法数据计算完成后要“看见”。可视化我用的是 Flask ECharts 这套轻量组合。Flask 作为后端接口层从 PostgreSQL 读取聚合结果ECharts 做前端图表渲染。可视化服务的 Dockerfile 关键部分FROM python:3.10-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt -i https://pypi.tuna.tsinghua.edu.cn/simple COPY app/ ./app/ EXPOSE 5000 CMD [gunicorn, -w, 4, -b, 0.0.0.0:5000, app:app]后端接口做实时数据拉取用 Flask 写一个返回 JSON 的接口app.route(/api/district/order/stat) def district_order_stat(): cursor pg_conn.cursor() cursor.execute( SELECT window_end, district, order_count, active_driver_count FROM district_order_stat ORDER BY window_end DESC LIMIT 100 ) rows cursor.fetchall() return jsonify({ timeline: [row[0].strftime(%H:%M) for row in reversed(rows)], districts: list({row[1] for row in rows}), series: [_build_series(rows)] })前端用 ECharts 的折线图展示每分钟订单趋势。关键点在于前端不能粗暴地每秒轮询否则后端数据库扛不住。我用的方案是前 30 秒轮询一次加上图表的数据追加模式这样既保证数据新鲜度又把后端压力控制住。关于 Flask 可视化部署我踩过一个比较值得记录的坑直接使用 Flask 内置的开发服务器是单线程的并发一高就出现请求排队、接口超时。后来换成了 Gunicorn 多 worker并发能力提升明显。如果容器内存允许建议 worker 数量设为 2 到 4 个太多 worker 反而会因为 Python GIL 的限制而没有明显收益还白占内存。4. 数据质量与平台稳定性保障机制4.1 数据质量检查清洗前后的六项校验规则数据工程里有一句流传很广的话垃圾进垃圾出。如果清洗环节没有质量把控数据仓库里的结论就没人敢信。我在管道里部署了一套独立的数据质量检查框架核心是六项校验规则规则检查内容执行时机失败处理非空检查关键字段是否有 NULL清洗后阻断下游报警通知唯一性检查订单 ID 是否重复清洗后记录重复率超阈值预警值域检查金额是否为负数、GPS 坐标是否越界清洗后剔除脏数据记录条数完整性检查当日数据总条数 vs 上游源数据条数入仓后偏差超 1% 时阻断时效性检查数据时间戳是否在合理窗口内入库时标记延迟数据单独分区一致性检查汇总表与明细表数值交叉核对建模后不一致时对比差异明细以订单日志清洗任务为例我在 Spark 作业里加了三个 Accumulator 计数器分别统计原始记录数、清洗后记录数、异常记录数。任务结束之后把这些信息打到日志里由 Athos自研监控捕获并判断是否符合预期。如果某个时间窗口的异常率突然从正常状态的 0.5% 飙到 5%说明上游业务线可能出了状况比如前端改了日志格式、数据库表结构变更质量检查框架就会触发告警并把问题作业挂起——宁可今天不出这张报表也不能让错误数据进入数仓和下游决策系统。4.2 数据血缘追踪找到这个表被谁在用数据平台规模大到一定程度数据血缘就变成了刚需。一个平台上有几千张表、几百个任务改了一个字段的语义到底会影响多少下游任务没有血缘追踪这个问题只能靠吼“谁用了这个表后端改一下你们注意下。”我们的实现方案比较落地基于 SQL 解析器对任务代码做静态解析。DolphinScheduler 调度平台记录每个工作流的任务 SQL我写了一个解析器去提取 SQL 里的 INSERT INTO 和 FROM 子句形成表级别的依赖关系最终写入 Neo4j 图数据库前端做成血缘可视化界面。实践效果如何有一个真实案例业务团队要下线一个旧的订单明细表用血缘工具一查发现下游有 27 个任务直接或间接依赖这张表。如果没有血缘工具这 27 个任务全部会在无人知晓的情况下静默失败尤其是那些跑在凌晨没有实时告警的批处理任务这种问题会被埋藏到早晨业务报告救不回来的境地。4.3 监控告警体系从指标采集到告警触达的完整链路云原生架构的可观测性不像物理机时代那样靠登录服务器敲命令而是全面转向指标 日志 追踪三位一体的体系。指标维度我用 Prometheus 的 kube-state-metrics 和 node-exporter 采集集群和节点的 CPU、内存、磁盘、网络指标cAdvisor 采集容器级别的指标Spark 和 Flink 都提供了原生的 Prometheus 指标暴露端点直接配置 scrape job 抓取即可。以 Spark 任务监控为例最需要关注的是这几个指标spark_executor_total_cores/spark_executor_used_cores看 Executor 是否真正忙起来了。我在调优参数时发现任务跑得慢的很大一部分原因是资源没有用满Executor 拿到 2 核但实际只用了 0.5 核白白浪费。spark_stage_failed_tasks/spark_task_failed失败任务数持续增长预示数据倾斜或者资源竞争。spark_executor_max_memory/spark_memory_used内存水位超过 80% 就要警惕 GC 频繁。告警触发后的通知渠道我统一收敛到企业微信机器人。通过 Alertmanager 配置路由规则将不同严重级别的告警发送到对应群组P0 严重集群宕机、数据管道中断超过 30 分钟同时走电话、短信、机器人。P1 警告任务失败率超过 10%、节点磁盘使用率超 85%机器人推送并 值班人。P2 提示数据质量校验小范围异常、任务运行时长偏离基线机器人推送第二天 Review 处理。这套监控系统大概上线跑了不到两个星期就帮我抓到了一个非常隐蔽的问题某个 Spark 任务之前一直运行 40 分钟某天开始逐渐增加到 90 分钟后触发了告警。排查后发现问题出在存储侧某个 PVC 的 IOPS 被打了上限导致 Spark 的 shuffle 写盘变慢。如果没有监控这个任务会继续每天慢下去直到业务方开口报障。5. 常见问题与排查技巧实录5.1 任务一直处于 Pending 状态的排查思路这是 Spark on K8s 场景最多见的问题新任务提交以后Driver Pod 创建成功但 Executor 迟迟起不来一直 Pending。我的排查顺序第一看资源。kubectl describe pod executor-pod-name如果 Events 里出现Insufficient memory或者Insufficient cpu就说明集群资源不够需要检查节点资源情况。我这里踩过坑调度器默认的resource配置和 Pod 实际 request 不一致明明节点还有内存但因为 Pod request 设置太大而无法调度。解决办法是调整 PodTemplate 里的资源申请值。第二看污点。K8s 默认会给节点打一些 taint比如node.kubernetes.io/disk-pressure如果 Executor Pod 没有设置对应的 toleration调度器就直接跳过这个节点。第三看镜像拉取。ImagePullBackOff或者ErrImagePull事件说明镜像仓库地址写错或者镜像 tag 不存在。检查私有仓库的认证配置imagePullSecret是否在serviceAccount里配置正确这一点很容易被遗忘。第四看配额。如果命名空间设置了ResourceQuotaPod 的请求资源超过了额度也会一直 Pending。通过kubectl describe resourcequota -n namespace查看配额使用情况。第五查调度器日志。kubectl logs -n kube-system kube-scheduler-xxxx --tail500里会有具体的调度失败原因。很多排不出来的“不可见”问题都是在这里找到真相的。5.2 数据倾斜导致 Spark 任务跑不完的优化记录有一张订单表的清洗任务按订单日期分区每个分区里有大量按城市字段做 group by 的操作。某天的任务运行时间从正常的 30 分钟涨到了 2 小时打开了 Spark UI 一看大多数 Task 在几十秒内就跑完了但最后一个 Task 跑了 1 个多小时——这是典型的数据倾斜。倾斜的原因很直接有个超大城市的订单量是其他城市的几十倍group by 时数据全部集中到一个任务里。对策是加一层随机前缀打散// 原始写法 df.groupBy(city) .agg( count(order_id).as(order_count), countDistinct(driver_id).as(driver_count) ) // 打散预处理 val saltedDf df.withColumn(city_salt, concat(col(city), lit(_), lit(rand() * 10))) saltedDf.groupBy(city_salt, city) .agg( count(order_id).as(order_count), countDistinct(driver_id).as(driver_count) ) .groupBy(city) .agg( sum(order_count).as(order_count), sum(driver_count).as(driver_count) )把倾斜组的压力分成 10 份并发处理再在最后一步 merge。改造之后同一个任务从 2 小时降到 32 分钟。另外还有一个容易被忽略的问题countDistinct(driver_id)在数据量大时非常耗时因为要维护一个去重集合。我后来把每小时去重改成用近似去重approx_count_distinct性能提升了一个数量级精确度误差完全在业务可接受范围内相对误差约 2%。对实时大屏这种对精确度要求不是特别高的场景近似算法是很值得投入的优化手段。5.3 对象存储读取性能慢的终极解法存储计算分离之后最常见的抱怨就是“以前 HDFS 上跑得挺快换到 MinIO 之后变慢了很多。”这个问题背后的本质也很清楚对象存储的文件访问走 HTTP 协议每读一次文件都要建立连接、鉴权、传输如果大量小文件并发访问网络开销会直接把性能拖垮。解决措施按效果排序输入侧做文件合并清洗任务把小文件合并成 128MB 或 256MB 的大文件。消除了大量的元数据请求读取吞吐量能提升一倍以上。用列式存储格式Parquet 或 ORC 格式配合谓词下推只读取需要的列和行组IO 量大幅缩小。缓存热点数据Spark 作业里对反复使用的 DataFrame 做.cache()或者在应用层开发一个本地缓存层将频繁访问的维度表放进内存。调整读取并发参数spark.sql.files.maxPartitionBytes从默认的 128MB 调大到 256MB让每个任务处理更多数据减少任务数量、降低调度开销。数据预分区按日期做目录分区查询只扫对应分区的文件避免全表扫描。做了以上五件事之后我这边对象存储读取的平均性能已经达到 HDFS 本地读取的 80% 左右剩下的差距由分布式存储的水平扩展能力来弥补——不够了加节点就行不用像 HDFS 那样调 balancing 和副本策略。5.4 常见问题速查表症状可能原因快速检查点Executor 一直 Pending资源不足 / 污点 / 镜像拉取失败kubectl describe pod看 Events任务运行缓慢数据倾斜 / 文件太小 / 资源未用满Spark UI 看 Task 分布和 Executor 利用率Flink 频繁重启内存不足 / RocksDB 本地空间不够看 JobManager 日志和 TaskManager 内存指标Kafka 消费积压持续增长下游计算能力不足 / 单分区热点看 consumer lag 和各分区 offset 分布写入 PostgreSQL 报连接池耗尽连接池太小 / 写入频率过高看连接池 max 值和当前 active 数告警风暴监控阈值设置不合理 / 下游连坐调整阈值并配置告警抑制规则6. 个人实操心得与后续演进方向6.1 镜像版本管理不做好这件事后期一定痛苦整个云原生架构里我最大的体会是镜像版本管理。一份镜像打错了 tag推到仓库覆盖了线上版本会让全平台所有任务一夜之间跑出完全不可预期的结果。我们在后期收紧了流程每个镜像必须有 git commit hash 作为 tag比如spark-3.4.0-ab12cd4。同时镜像一旦发布到生产环境使用的 tag禁止覆盖更新只允许新增。部署时在values.yaml或 PodTemplate 中指定精确的 tag用 Helm 的版本管理来追踪配置变化。这样出了问题能快速定位到是哪一次变更引入的——git diff一眼看清。6.2 成本治理的一点经验云原生弹性好但若不加控制成本会像漏水的桶一样无声无息地流走。我做了两件事来控制第一预算大盘。按部门/项目维度统计 CPU 和内存消耗输出每日成本报告让各业务线看到自己消耗的资源对应的成本。数据一透明浪费自然减少。第二闲置资源回收。通过 CronJob 每天凌晨扫描集群中全天 CPU 使用率低于 5% 的 Workload自动缩容到 0白天再按需拉起。这一步实测能节约 10%-20% 的计算资源成本。6.3 架构演进的一个方向从批流一体到数据湖仓做完这整套系统我一直在思考下一步。目前离线链路和实时链路是分开的离线走 Spark 批处理实时走 Flink 流处理。这带来的问题是同一份业务数据被计算了两遍而且批流之间对同一份数据的口径可能不一致。后续迭代方向是引入数据湖仓技术比如 Apache Iceberg 或 Paimon用流批一体的方式管理表格。核心思路是实时任务写数据湖的增量分区离线任务读数据湖的全量分区同一份数据存储层统一、查询口径统一。目前我们已经在部分业务场景试点了 Iceberg流批一体在保存点恢复、Schema 演进方面的能力确实让数据管道的维护成本又降了一个台阶。这个方向我也还在实践过程中等跑得足够充分了再单独写一篇完整分享。