云原生 AI 与大数据:Kubeflow、Spark、Flink 与特征平台

2022-02-11T10:00:00+08:00 | 13分钟阅读 | 更新于 2022-02-11T10:00:00+08:00

@

云原生 AI 与大数据:从「抢显卡」到「存算分离」

云原生不只是把 Web 服务塞进容器。AI 训练要抢 GPU,Spark 要弹性 Executor,Flink 要** Exactly-Once 实时特征**,推理服务要按 GPU 指标扩缩,数据湖要管元数据。这一章把第 12 章拆成五条线,每条线先用一句白话建直觉,再给可运行的配置 / 代码,最后点出面试怎么答。

  • 训练线:GPU 是稀缺停车位,Argo + GPU Pooling 让你「先到先停、临时车可抢占」。
  • Spark 线:Executor 是临时工,request.cores 决定他占几个「核工位」,Dynamic Allocation 决定活少就裁员。
  • 特征线:离线和在线特征像「两份菜谱」,Feast 保证它们一致;热点 Key 像爆款商品,要分桶打散。
  • 推理线:Triton 像连锁快餐店,Warmup 是开门前先热好油锅,KEDA 按 GPU 利用率加开店。
  • 数据湖线:Iceberg 是带版本的书架,HMS 是图书管理员,孤儿文件是没人认领的废书。

一、训练编排与 GPU 弹性(12.1)

直觉:GPU 是稀缺车位,训练任务是来停车的车

GPU 节点贵且少。多个训练任务抢同一池子 GPU 时,要让短任务能抢占长任务(就像临停可以客气地请长时间占位的车挪一挪),还要能动态加 worker(数据并行训练,数据多了就多开几台机器一起算)。

graph TD
  WF[Argo Workflow:train] --> P1[Pod worker-0 GPU]
  WF --> P2[Pod worker-1 GPU]
  WF --> P3[Pod worker-2 GPU]
  Pool[GPU Pooling:共享 GPU 节点池] -. 弹性抢占 .-> P1
  Pool -. 动态扩缩 .-> P2
  Topo[Pod 拓扑分布约束] -. 数据倾斜时重调度 .-> P3

12.1.1 Argo Workflows + GPU Pooling 弹性抢占

GPU Pooling 通过把 GPU 节点放进一个共享节点池,并用 PriorityClass / 抢占式调度让高优任务抢低优任务的 GPU。Argo Workflow 的 template 里用 nodeSelector 指向 GPU 池,用 priorityClassName 声明抢占级别。

# argo-gpu-train.yaml —— 训练 Workflow,可抢占低优任务
apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
  name: gpu-train
spec:
  entrypoint: train
  templates:
    - name: train
      dag:
        tasks:
          - name: worker
            template: torch-train
            withSequence:               # 动态起 N 个 worker
              count: "{{workflow.parameters.workers}}"
    - name: torch-train
      nodeSelector:
        cloud.google.com/gke-accelerator: nvidia-tesla-t4   # 指定 GPU 池
      priorityClassName: gpu-preemptible                    # 高优可抢占
      container:
        image: pytorch:2.1
        resources:
          limits:
            nvidia.com/gpu: 1
        command: ["python", "train.py"]

12.1.2 KFP SDK 数据并行训练并动态调整 worker 数(Python)

Kubeflow Pipelines (KFP) 用 Python 定义 Pipeline。下面这段用 create_component 抽象训练步骤,并通过 withParametersworker 数作为运行时参数传入,实现「数据并行 + 动态 worker」。

from kfp import dsl
from kfp.dsl import component, pipeline

# 1) 把训练脚本包成一个组件
@component(base_image="pytorch:2.1")
def train_worker(epochs: int, worker_id: int, world_size: int):
    import os
    # 数据并行:每个 worker 用不同的分片数据,world_size 即总 worker 数
    os.environ["WORLD_SIZE"] = str(world_size)
    os.environ["RANK"] = str(worker_id)
    print(f"worker {worker_id}/{world_size} training...")

# 2) Pipeline:动态起 world_size 个 worker(值可在触发时改)
@pipeline(name="data-parallel-train")
def train_pipeline(epochs: int = 5, world_size: int = 4):
    for i in range(world_size):
        train_worker(epochs=epochs, worker_id=i, world_size=world_size)

# 3) 编译 & 提交(动态改 worker 数只需改 world_size 参数)
if __name__ == "__main__":
    from kfp import compiler
    compiler.Compiler().compile(train_pipeline, "train.yaml")

面试要点:world_size 是数据并行的核心——它等于「把数据集切成几份同时算」。动态调整 worker 数本质是改 world_size 并重调度,需要训练框架(PyTorch DDP / TF 的 tf.distribute)支持弹性 rendezvous

12.1.3 Pipeline 数据倾斜时用 Pod 拓扑分布约束重新调度

数据倾斜 = 某些 worker 分到的数据特别多、算得特别慢,整体被拖垮。Kubernetes 的 拓扑分布约束(topologySpreadConstraints) 可以把 worker 均匀铺在不同节点/可用区,避免「都挤在一台机器上抢资源」,必要时结合反亲和性重调度。

# topology-spread.yaml —— 训练 worker 均匀打散,缓解倾斜
topologySpreadConstraints:
  - maxSkew: 1                 # 各节点 Pod 数差不超过 1
    topologyKey: kubernetes.io/hostname
    whenUnsatisfiable: ScheduleAnyway   # 实在不行也先跑,但尽量匀
    labelSelector:
      matchLabels:
        app: train-worker
# 倾斜严重时可加反亲和,强制分散
affinity:
  podAntiAffinity:
    requiredDuringSchedulingIgnoredDuringExecution:
      - labelSelector:
          matchLabels: { app: train-worker }
        topologyKey: kubernetes.io/hostname

二、Spark on Kubernetes(12.2)

直觉:Executor 是临时工,核是工位,空闲就裁员

Spark 在 K8s 上每个 Executor 是一个 Pod。两个核心调优点:别让一个 Executor 占的「核工位」数踩到超线程导致互相干扰;活少的时候让 Dynamic Allocation 把闲着的 Executor 裁掉省资源。

graph LR
  Driver[Spark Driver Pod] -->|申请 Executor| API[K8s API]
  API --> E1[Executor Pod:2 核]
  API --> E2[Executor Pod:2 核]
  API --> E3[Executor Pod:2 核]
  DA[Dynamic Allocation] -. 空闲回收 .-> E3
  ST[Shuffle Tracking] -. 保留 shuffle 数据 .-> E1

12.2.1 spark.kubernetes.executor.request.cores 避免超线程干扰

一个物理核通常有 2 个超线程。如果 executor.cores=2request.cores 没限制,K8s 可能把两个超线程算成「2 核」调度,实际两个 Executor 挤在同一物理核上互相抢。设置 request.cores 让 K8s 按真实物理核请求,避免超线程干扰。

# spark-k8s-cores.yaml —— spark-submit 等价配置片段
apiVersion: v1
kind: ConfigMap
metadata:
  name: spark-conf
data:
  spark-defaults.conf: |
    # 每个 Executor 用 2 个逻辑核
    spark.executor.cores                     2
    # 关键:向 K8s 请求的核数,避免把超线程当整核调度
    spark.kubernetes.executor.request.cores  2
    # 限制单节点最多跑几个 Executor,进一步隔离
    spark.kubernetes.executor.limit.cores    4    

面试要点:spark.executor.cores 是 Spark 内部并行度;spark.kubernetes.executor.request.cores给 K8s 调度器的请求值。两者不一致时,调度按后者,计算按前者——不一致就会超线程撞车。

12.2.2 Dynamic Allocation + Shuffle Tracking 减少空闲 Executor

Dynamic Allocation 让 Executor 数量随任务量弹性变化;但 Executor 被回收后,它产生的 shuffle 中间数据会丢,下游任务就得重算。开启 Shuffle Tracking 后,K8s 上的 shuffle 数据由外部跟踪服务(如 Spark 的 shuffle 服务 / RSS)保留,Executor 可安心回收。

# spark-dynamic.yaml —— ConfigMap 中的动态分配配置
apiVersion: v1
kind: ConfigMap
metadata:
  name: spark-dynamic
data:
  spark-defaults.conf: |
    spark.dynamicAllocation.enabled          true
    spark.dynamicAllocation.minExecutors    2
    spark.dynamicAllocation.maxExecutors    20
    spark.dynamicAllocation.executorIdleTimeout 60s   # 空闲 60s 回收
    # 关键:开启 shuffle 跟踪,回收 Executor 不丢中间数据
    spark.dynamicAllocation.shuffleTracking.enabled true
    # 配合外部 shuffle 服务(如 Kubernetes shuffle service / RSS)
    spark.shuffle.service.enabled            true    

12.2.3 S3A 提交出现 FileNotFoundException 时启用 EMRFS 一致性视图

Spark 用 S3A 写数据到 S3 时,S3 是最终一致性的:文件刚写完,Listing 可能还看不到,于是 commit 阶段报 FileNotFoundException。AWS 的 EMRFS 一致性视图(Consistency View) 用一张 DynamoDB 表记录「已成功写入的对象」,让 Spark 提交时以这张表为准,绕过 S3 一致性延迟。

# spark-emrfs.yaml —— 启用 EMRFS 一致性视图
apiVersion: v1
kind: ConfigMap
metadata:
  name: spark-emrfs
data:
  spark-defaults.conf: |
    # 开启 EMRFS 一致性视图
    spark.hadoop.fs.s3a.consistent.status  true
    # 一致性元数据存到 DynamoDB 表
    spark.hadoop.fs.s3a.consistent.table   emrfs-consolidated
    # 重试与重试间隔,进一步兜住瞬时不可见
    spark.hadoop.fs.s3a.consistent.retry.count   5
    spark.hadoop.fs.s3a.consistent.retry.interval 200ms    

面试要点:根因是「S3 最终一致性 vs Spark commit 的强一致假设」。除了 EMRFS,也可换用 S3 的目录提交器(directory committer / manifest committer) 或对象存储本身的强一致(如 OSS、GCS)来缓解。


三、特征平台与实时计算(12.3)

直觉:离在线特征是「两份菜谱」,热点 Key 是爆款商品

特征平台要把离线算好的特征(批处理,T+1)和在线实时特征(毫秒级)对齐,否则训练和推理用的特征不一致,模型就「学的是一套、用的是另一套」。Flink 实时算特征写入 Redis 时要保证不重不漏(Exactly-Once)。热点 Key 就像秒杀爆款,所有请求砸同一个 Redis 分片,要按一致性哈希分桶打散

graph LR
  Off[离线特征批计算] --> FS[(Feast 特征仓库)]
  RT[Flink 实时计算] --> Redis[(在线特征 Redis)]
  FS --> Serv[推理服务]
  Redis --> Serv
  Feast[Feast 对齐离/在线] -. 同一套特征定义 .-> FS
  Redis -. Consistent Hash 分桶 .-> H[热点 Key 打散]

12.3.1 Feast on Kubernetes 保证离在线特征一致性

Feast 的核心思想是:离线和在线特征共用同一份特征定义(Feature View)和实体键。离线存在仓库(如 Iceberg/Parquet),在线存在 Redis,但「特征怎么算、主键是什么」只定义一次,从而保证两边口径一致。下面是用 Feast SDK 注册一个特征视图并物化到在线的例子。

from feast import FeatureStore, FeatureView, Entity, ValueType, FileSource, RedisSource
from feast import Field
from datetime import timedelta

# 1) 定义实体(特征挂在哪个主键上)
user = Entity(name="user_id", value_type=ValueType.INT32)

# 2) 离线源(批计算产出的 Parquet)
offline_source = FileSource(
    path="s3a://features/offline/user_features.parquet",
    timestamp_field="event_ts",
)

# 3) 在线源(Redis),与离线用同一个特征定义
online_source = RedisSource(
    connection_string="redis://redis.feat.svc:6379",
    key="user_id",
)

# 4) 特征视图:离/在线共用,口径天然一致
view = FeatureView(
    name="user_features",
    entities=[user],
    ttl=timedelta(days=1),
    schema=[Field(name="avg_price", dtype=ValueType.FLOAT)],
    online=True,
    source=offline_source,
)

# 5) 注册并物化到在线存储
store = FeatureStore(repo_path="./feast_repo")
store.apply([user, view])
store.materialize_incremental(end_date="2024-01-01")   # 增量刷到 Redis

Flink 通过 checkpoint + 两阶段提交(2PC) 实现 Exactly-Once。写入 Redis 时用支持事务的 connector,并在 DDL 里声明 sink.semantics = exactly-once。下面把订单流和用户特征流 Join,产出实时特征写入 Redis。

-- 1) 订单流(Kafka 源,带水位线保证乱序处理)
CREATE TABLE orders (
  user_id INT,
  amount  DOUBLE,
  ts      TIMESTAMP(3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic'     = 'orders',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format'    = 'json'
);

-- 2) 用户特征流(另一路 Kafka)
CREATE TABLE user_feat (
  user_id   INT,
  avg_price DOUBLE,
  ts        TIMESTAMP(3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH ('connector'='kafka','topic'='user_feat','properties.bootstrap.servers'='kafka:9092','format'='json');

-- 3) 实时特征 Join 结果写入 Redis,声明 exactly-once 语义
CREATE TABLE feat_out (
  user_id   INT,
  real_avg  DOUBLE,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'connector'          = 'redis',
  'mode'               = 'single',
  'host'               = 'redis.feat.svc',
  'port'               = '6379',
  'sink.semantics'    = 'exactly-once',     -- 关键:依赖 checkpoint + 2PC
  'sink.flush.interval' = '1000'
);

-- 4) 双流 Join 写入
INSERT INTO feat_out
SELECT o.user_id, (o.amount + u.avg_price) / 2.0
FROM orders o
JOIN user_feat u ON o.user_id = u.user_id;

面试要点:Exactly-Once 不是「每条只处理一次」那么玄,而是 checkpoint 对齐 + 下游幂等/事务写入。Redis connector 的 exactly-once 靠 checkpoint 屏障和写入幂等键实现;若 Redis 不支持事务,要退而用「唯一键 + upsert」保证不重。

12.3.3 特征存储热点 Key 基于 Consistent Hash 分桶打散

当某个用户/商品是爆款,所有特征请求都打同一个 Redis 分片,分片被打挂。解决:对 Key 做一致性哈希分桶——在原始 Key 后拼一个桶号(如 user:123#bucket3),把热点摊到多个分片。下面用 Python 演示分桶逻辑。

import hashlib

def shard_key(raw_key: str, bucket_count: int = 16) -> str:
    """对热点 key 做一致性哈希分桶,把压力摊到多个分片"""
    # 1) 用原始 key 算一个哈希,映射到 [0, bucket_count)
    h = int(hashlib.md5(raw_key.encode()).hexdigest(), 16)
    bucket = h % bucket_count
    # 2) 拼桶号,写入/读取都用这个分桶 key,保证同 key 落到同桶
    return f"{raw_key}#b{bucket}"

# 写入:热点用户 123 被拆到 16 个桶之一
for _ in range(1000):
    k = shard_key("user:123")      # 结果稳定,如 "user:123#b5"
    redis.set(k, feature_value)

面试要点:一致性哈希保证「同 key 永远落同桶」,分桶数要选得够大以摊均热点,又不能太大导致小 key 也碎片化。配合 Redis Cluster 的 slot 分布效果更佳。


四、Triton Inference Server(12.4)

直觉:Triton 是连锁快餐店,Warmup 是开门前热油锅

Triton 是 NVIDIA 的推理服务器,一个 Pod 能同时托管多个模型。首请求往往特别慢(要加载模型、分配显存、编译 kernel),这叫「冷启动延迟」。Model Warmup 就是开门营业前先把油锅烧热——启动时主动跑几遍 dummy 请求,让首波真实流量不再挨冷。

graph TD
  Req[推理请求] --> Triton[Triton Pod]
  Triton --> Warm[Model Warmup:启动时预热]
  Triton --> Infer[GPU 推理]
  KEDA[KEDA + GPU Metrics] -. 显存/利用率高 .-> Scale[扩容 Triton Pod]
  Offload[Model Ensemble CPU Offload] -. 显存不足 .-> CPU[部分层放 CPU]

12.4.1 Triton Model Warmup 避免首次延迟

在模型的 config.pbtxt 里加 model_warmup 段,指定预热时用的数据形状和批次,Triton 启动时自动跑。

# config.pbtxt —— ResNet50 的预热配置
name: "resnet50"
platform: "tensorrt_plan"
max_batch_size: 8
# 关键:启动时预热,避免首请求冷启动
model_warmup [
  {
    name: "warmup_bs8"
    batch_size: 8
    # 用随机数据跑 2 次,填满显存分配与 kernel 编译路径
    inputs: { key: "input", value: { data_type: TYPE_FP32, dims: [3,224,224], zero_data: true } }
    repeat: 2
  }
]

12.4.2 KEDA + GPU Metrics 自动扩容 Triton Pod 的 ScaledObject

用 KEDA 的 Prometheus 触发器读 GPU 利用率(来自 DCGM exporter),利用率超阈值就加 Triton Pod。

# triton-scaledobject.yaml —— 基于 GPU 利用率弹性扩容
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: triton-scaler
  namespace: inference
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: triton
  minReplicaCount: 1
  maxReplicaCount: 10
  triggers:
    - type: prometheus
      metadata:
        serverAddress: http://prometheus.monitoring.svc:9090
        metricName: gpu_util
        # 查询 DCGM 的 GPU 利用率
        query: 'DCGM_FI_DEV_GPU_UTIL{namespace="inference"}'
        threshold: "70"          # 利用率超 70% 触发扩容
        activationThreshold: "40"

12.4.3 显存不足时启用 Model Ensemble CPU Offload 保证 P99<200ms

当多个模型同驻一块 GPU 显存不够,可把 Ensemble 里预处理/后处理或部分层 offload 到 CPU,只把核心推理留 GPU,从而显存够用且 P99 不超标。Triton 通过 Ensemble 把不同 backend 串联,CPU 部分用 onnxruntime/python backend。

# ensemble_config.pbtxt —— 把预处理 offload  CPU
name: "ens_resnet"
platform: "ensemble"
# 组成:cpu_pre(CPU) -> gpu_infer(GPU) -> cpu_post(CPU)
ensemble_scheduling {
  step [
    { model_name: "preprocess_cpu"; model_version: -1; input_map { key: "RAW"; value: "RAW" }; output_map { key: "PREP"; value: "PREP" } },
    { model_name: "resnet50_gpu";  model_version: -1; input_map { key: "input"; value: "PREP" }; output_map { key: "output"; value: "LOGIT" } },
    { model_name: "postprocess_cpu"; model_version: -1; input_map { key: "LOGIT"; value: "LOGIT" }; output_map { key: "RESULT"; value: "RESULT" } }
  ]
}

面试要点:offload 的代价是 CPU-GPU 数据传输,但要算总账——显存够 → 不 OOM → 不降级 → P99 稳。配合 12.4.1 的 Warmup 和 12.4.2 的扩容,三者共同兜住延迟 SLO。


五、数据湖 Iceberg on K8s(12.5)

直觉:Iceberg 是带版本的书架,HMS 是管理员,孤儿文件是废书

Iceberg 是表格式(table format),给数据湖加上快照、隐藏分区、时间旅行。元数据存在 HMS(Hive Metastore)或兼容服务里;HMS Pool 是连接池,防止高并发时元数据请求把单点 HMS 打爆。快照清理会留下孤儿文件(没人引用的旧数据文件),要定期回收。

graph LR
  Spark[Spark on K8s] --> Iceberg[(Iceberg 表)]
  Iceberg --> HMS[HMS Pool 元数据连接池]
  Iceberg --> Data[Parquet/ORC 数据文件]
  Alluxio[Alluxio + Job Master] -. 批量回收孤儿文件 .-> Data
  Snap[快照自动清理] -. 过期快照 .-> Data

12.5.1 Iceberg on K8s 设置 HMS Pool 避免元数据请求超时

Spark 连 HMS 时若每次新建连接,高并发下 HMS 单点易被拖垮超时。用 HMS Pool(连接池):通过 hive.metastore.client.pool.size 复用连接,并加超时与重试。

# iceberg-hms-pool.yaml —— Spark 访问 Iceberg 的 HMS 连接池配置
apiVersion: v1
kind: ConfigMap
metadata:
  name: iceberg-conf
data:
  spark-defaults.conf: |
    # Iceberg catalog 类型
    spark.sql.catalog.iceberg         iceberg
    spark.sql.catalog.iceberg.type    hive
    spark.sql.catalog.iceberg.uri     thrift://hms.metastore.svc:9083
    # 关键:HMS 客户端连接池,避免元数据请求排队超时
    spark.hadoop.hive.metastore.client.pool.size     8
    spark.hadoop.hive.metastore.client.connect.retry.limit 5
    spark.hadoop.hive.metastore.client.socket.timeout 30s    

12.5.2 Spark + Iceberg 隐藏分区并自动清理快照的 SQL

隐藏分区(Hidden Partitioning):分区字段由 Iceberg 根据数据自动派生(如按 tsday 分区),写入/查询时用户不必手动指定分区值,避免「写错分区」的经典 bug。下面创建带隐藏分区的表,并配置快照自动清理

-- 1) 建表:按事件时间的天做隐藏分区,用户无需手写分区列
CREATE TABLE iceberg.db.user_events (
  user_id    BIGINT,
  event_ts   TIMESTAMP,
  payload    STRING
) USING iceberg
PARTITIONED BY (days(event_ts));   -- 隐藏分区:按天,查询自动路由

-- 2) 写入(不用写 PARTITION,Iceberg 自动算)
INSERT INTO iceberg.db.user_events VALUES (1, TIMESTAMP '2024-01-01 10:00:00', 'x');

-- 3) 配置快照生命周期:保留 7 天,自动清过期快照与孤儿数据
ALTER TABLE iceberg.db.user_events SET TBLPROPERTIES (
  'snapshot.retain-last' = '10',
  'snapshot.expire.max-ref-age-ms' = '604800000',  -- 7 天
  'write.delete.orphan.files.enabled' = 'true'      -- 清孤儿文件
);

-- 4) 手动触发快照过期(也可交给定时 Spark Job)
CALL iceberg.system.expire_snapshots('db.user_events', TIMESTAMP '2024-01-01 00:00:00');

12.5.3 元数据孤儿文件基于 Alluxio + Job Master 批量回收

即使开了快照清理,残留的孤儿文件(历史快照引用过的 Parquet,但当前无引用)仍可能留在存储里占空间。Alluxio 的 Job Master 能跨节点批量扫描并 DELETE 这些孤儿,避免单点 Spark Driver 内存扛不住。

# alluxio-orphan-clean.yaml —— 用 Alluxio Job Master 跑批量回收 Job
apiVersion: batch/v1
kind: Job
metadata:
  name: iceberg-orphan-clean
  namespace: data
spec:
  template:
    spec:
      restartPolicy: OnFailure
      containers:
        - name: cleaner
          image: alluxio/alluxio:3.1
          command:
            - /bin/sh
            - -c
            - |
              # 通过 Alluxio Job Master 提交孤儿文件回收任务
              # 扫描 Iceberg 元数据目录,删除无引用的数据文件
              alluxio job submit \
                --type IcebergOrphanClean \
                --path s3a://lake/db/user_events \
                --expired-before 2024-01-01              

面试要点:孤儿文件回收要先确认快照真的没人引用(时间旅行查询可能还指着旧快照),否则会误删正在被读的历史版本。生产上回收 Job 要加「仅删早于 N 天且不在任何有效快照中」的保险。


自测题与动手练习

  1. 概念题:Argo Workflows 配 GPU Pooling 抢占时,靠哪两个 K8s 机制实现「高优抢低优」?KFP 动态调 worker 数本质是改哪个变量?
  2. 代码题:写一段 Flink SQL,把订单流和用户特征流 Join 后 Exactly-Once 写入 Redis,指出保证不重不漏的三个关键点。
  3. 排错题:Spark 往 S3 提交报 FileNotFoundException,根因是什么?给出 EMRFS 一致性视图的核心配置项。
  4. 配置题:Triton 显存不足时,如何用 Model Ensemble 把预处理 offload 到 CPU?配合哪两个手段共同兜住 P99<200ms?
  5. 设计题:Iceberg 隐藏分区相比 Hive 传统分区的优势是什么?孤儿文件回收为什么不能无脑删,要加什么保险?

动手练习:在 minikube 上用 Helm 装 Spark Operator,提交一个带 request.cores 和 Dynamic Allocation 的 Pi 作业,观察 Executor Pod 的 requests.cpu 是否等于 executor.cores;再用 Feast 官方 docker-compose 起一套离在线特征仓库,注册一个特征视图并 materialize 到 Redis。

本章小结

  • 训练编排:GPU Pooling + PriorityClass 实现弹性抢占;KFP 用 world_size 参数化数据并行 worker;拓扑分布约束 + 反亲和缓解数据倾斜。
  • Spark on K8s:request.cores 把超线程当整核请求避免干扰;Dynamic Allocation + Shuffle Tracking 在省资源的同时不丢中间数据;S3 最终一致性导致的 FileNotFoundException 用 EMRFS 一致性视图兜住。
  • 特征平台:Feast 用同一份特征定义保证离在线一致;Flink SQL 靠 checkpoint + exactly-once connector 实现实时特征 Join;热点 Key 用一致性哈希分桶打散。
  • Triton:Model Warmup 消除冷启动;KEDA 读 GPU 指标弹性扩容;Ensemble CPU Offload 在显存不足时保 P99。
  • Iceberg:HMS Pool 防元数据超时;隐藏分区避免写错分区;Alluxio Job Master 批量回收孤儿文件,但须先确认快照无引用。
复习提示:
  • GPU Pooling 핵심:PriorityClass 로 우선순위 매기고 Preemptible 노드에서 short job 이 long job占有 GPU 를 temporarily preempt 할 수 있게 함.
  • Spark Dynamic Allocationexecutor.cores request 로 core requests 설정, shuffle tracking 으로 intermediate data 잃지 않으면서 자원 절약.
  • Flink Exactly-Once:checkpoint + two-phase commit connector(예: Kafka producer transaction) 로 정확히 한 번 처리 보장.
  • Triton Ensemble:모델을 여러 단계로 나누어 CPU/GPU 에 분산 처리,显存 부족 시 P99 지연 보호.
  • Feast feature store:离线/online feature 일관성 유지, hot key 는 consistency hashing 분산으로 처리.
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!