ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

ZenML 大数据处理实战:从内存优化、分块计算到分布式引擎与 LakeFS 数据版本化的渐进式扩展方案

ZenML 大数据处理实战:从内存优化、分块计算到分布式引擎与 LakeFS 数据版本化的渐进式扩展方案 ZenML 大数据处理实战从内存优化、分块计算到分布式引擎与 LakeFS 数据版本化的渐进式扩展方案【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml当数据集不断膨胀单机 pandas 工作流终将触及内存与 I/O 的天花板。本文基于 ZenML 官方教程 Handling big data系统讲解如何渐进式地扩展一条 ZenML 管道先用 Parquet 与采样优化内存内处理再过渡到分块out-of-core与数据仓库卸载进而把 Spark、Ray、Dask 等分布式引擎直接嵌入 step最后用 LakeFS 为 50 GB 的重数据建立外部版本层。读完本篇你将掌握一套按数据规模选型的完整方法论并能结合仓库中的 LakeFS 示例代码落地轻量引用 重型数据的管道设计。理解数据集规模阈值先选对武器再动手在引入任何具体技术之前先建立规模感。教程给出了三个通用阈值区间不同区间对应不同的必要手段数据规模典型量级推荐策略小型数据集几 GB 以内标准 pandas 内存内操作即可中型数据集几十 GB分块chunking或 out-of-core 处理大型数据集几百 GB 及以上必须引入分布式处理框架这三档不是非黑即白的实际项目中应从简单策略起步随数据增长逐步升级。ZenML 的架构允许你在同一条管道中混合使用这些策略——例如 step A 用分块读取 CSVstep B 把聚合结果交给 BigQuerystep C 再用 Ray 做并行预处理。内存内工作流优化几 GB 以内数据仍能装进内存但已经开始吃力时有三类低成本高收益的优化。1. 换用高效数据格式CSV → ParquetCSV 是无类型、无压缩的文本格式读取时需要全量解析与类型推断Parquet 则是列式、压缩、带 schema 的二进制格式对数值型工作负载的读写都有数量级优势。教程给出的 Dataset 类示例import pyarrow.parquet as pq class ParquetDataset(Dataset): def __init__(self, data_path: str): self.data_path data_path def read_data(self) - pd.DataFrame: return pq.read_table(self.data_path).to_pandas() def write_data(self, df: pd.DataFrame): table pa.Table.from_pandas(df) pq.write_table(table, self.data_path)注意这里通过pyarrow.parquet直接操作 Arrow Table 再转 pandaspq.read_table().to_pandas()是零拷贝路径避免了经过 CSV 字符串解析的额外开销。如果你要深入自定义 Dataset 类的写法可参考教程的自定义 Dataset 类一篇。2. 引入基础采样用 10% 的数据验证 100% 的逻辑在开发迭代阶段全量跑一轮 pipeline 往往不现实。给 Dataset 类加一个采样方法让探索性 step 只处理子集import random class SampleableDataset(Dataset): def sample_data(self, fraction: float 0.1) - pd.DataFrame: df self.read_data() return df.sample(fracfraction) step def analyze_sample(dataset: SampleableDataset) - Dict[str, float]: sample dataset.sample_data(fraction0.1) # 在样本上做分析 return {mean: sample[value].mean(), std: sample[value].std()}采样 step 返回的是轻量统计量而非 DataFrame这类小结果对象走 ZenML artifact store 几乎没有成本——这也是后文 LakeFS 模式要放大的同一个思路。3. 优化 pandas / NumPy 操作减少不必要的内存拷贝优先使用向量化 NumPy 运算而非 Python 层循环step def optimize_processing(df: pd.DataFrame) - pd.DataFrame: # 尽量使用 in-place 操作 df[new_column] df[column1] df[column2] # 用 NumPy 运算换取速度 df[mean_normalized] df[value] - np.mean(df[value]) return df这一层的优化空间通常能覆盖几 GB 以内的全部场景是成本最低的升级路径。Out-of-core 处理几十 GB当数据装不进内存核心思路是不要让完整数据同时存在于内存中。分块读取大 CSV在 Dataset 类中用生成器实现分块每次只处理一个chunk_size大小的片段class ChunkedCSVDataset(Dataset): def __init__(self, data_path: str, chunk_size: int 10000): self.data_path data_path self.chunk_size chunk_size def read_data(self): for chunk in pd.read_csv(self.data_path, chunksizeself.chunk_size): yield chunk step def process_chunked_csv(dataset: ChunkedCSVDataset) - pd.DataFrame: processed_chunks [] for chunk in dataset.read_data(): processed_chunks.append(process_chunk(chunk)) return pd.concat(processed_chunks) def process_chunk(chunk: pd.DataFrame) - pd.DataFrame: # 在这里处理每个 chunk return chunkpd.read_csv(..., chunksizeN)返回的迭代器每次只把 N 行载入内存。需要注意的边界如果 step 最终要pd.concat所有处理结果则输出仍需整体装进内存——分块适合流式聚合、流式过滤这类可增量完成的计算若必须产出全量结果应考虑把结果直接写回 Parquet 或对象存储而非 DataFrame。把重 SQL 推给数据仓库对于聚合类计算本地分块往往不如直接把 SQL 交给云数据仓库如 Google BigQuery的分布式引擎step def process_big_query_data(dataset: BigQueryDataset) - BigQueryDataset: client bigquery.Client() query f SELECT column1, AVG(column2) as avg_column2 FROM {dataset.table_id} GROUP BY column1 result_table_id f{dataset.project}.{dataset.dataset}.processed_data job_config bigquery.QueryJobConfig(destinationresult_table_id) query_job client.query(query, job_configjob_config) query_job.result() # 等待作业完成 return BigQueryDataset(table_idresult_table_id)这个写法的关键点通过bigquery.QueryJobConfig(destination...)把聚合结果直接落到 BigQuery 的新表结果不经过本地内存step 只返回一个指向新表的BigQueryDataset引用对象。这与后文 LakeFS 的传递引用而非数据是完全一致的模式。接入分布式计算引擎几百 GB面对海量数据可以把 Spark、Ray、Dask 等引擎直接用进 ZenML step。需要说明的是ZenML 没有为这些框架提供内置集成组件但 step 函数内部就是普通 Python 代码因此可以直接初始化并使用这些框架——管道的编排、重试、版本化仍由 ZenML 承担。接入 Apache Sparkfrom pyspark.sql import SparkSession from zenml import step, pipeline step def process_with_spark(input_data: str) - None: # 初始化 Spark spark SparkSession.builder.appName(ZenMLSparkStep).getOrCreate() # 读取数据 df spark.read.format(csv).option(header, true).load(input_data) # 用 Spark 处理 result df.groupBy(column1).agg({column2: mean}) # 写出结果 result.write.csv(output_path, headerTrue, modeoverwrite) # 停止 Spark session spark.stop() pipeline def spark_pipeline(input_data: str): process_with_spark(input_data) # 运行管道 spark_pipeline(input_datapath/to/your/data.csv)前提是执行环境本地或远端容器中已安装 Spark 及其依赖。接入 Rayimport ray from zenml import step, pipeline step def process_with_ray(input_data: str) - None: ray.init() ray.remote def process_partition(partition): # 处理数据的一个分区 return processed_partition # 加载并切分数据 data load_data(input_data) partitions split_data(data) # 把处理任务分发到 Ray 集群 results ray.get([process_partition.remote(part) for part in partitions]) # 合并并保存结果 combined_results combine_results(results) save_results(combined_results, output_path) ray.shutdown() pipeline def ray_pipeline(input_data: str): process_with_ray(input_data) # 运行管道 ray_pipeline(input_datapath/to/your/data.csv)与 Spark 同理需要保证运行环境安装了 Ray。Ray 的分区 remote 任务模型对 pandas 生态非常友好适合把已有的单机 DataFrame 代码改造成并行版本。接入 Dask附自定义 MaterializerDask 提供与 pandas API 高度兼容的并行 DataFramedask.dataframe。由于 ZenML 内置 materializer 不认识dd.DataFrame教程展示了自定义 Materializer 的标准写法——继承 BaseMaterializer声明关联类型与保存/加载逻辑from zenml import step, pipeline import dask.dataframe as dd from zenml.materializers.base_materializer import BaseMaterializer import os class DaskDataFrameMaterializer(BaseMaterializer): ASSOCIATED_TYPES (dd.DataFrame,) ASSOCIATED_ARTIFACT_TYPE dask_dataframe def load(self, data_type): return dd.read_parquet(os.path.join(self.uri, data.parquet)) def save(self, data): data.to_parquet(os.path.join(self.uri, data.parquet)) step(output_materializersDaskDataFrameMaterializer) def create_dask_dataframe(): df dd.from_pandas(pd.DataFrame({A: range(1000), B: range(1000, 2000)}), npartitions4) return df step def process_dask_dataframe(df: dd.DataFrame) - dd.DataFrame: result df.map_partitions(lambda x: x ** 2) return result step def compute_result(df: dd.DataFrame) - pd.DataFrame: return df.compute() pipeline def dask_pipeline(): df create_dask_dataframe() processed process_dask_dataframe(df) result compute_result(processed) # 运行管道 dask_pipeline()这里有两个值得注意的设计细节Materializer 负责把计算图物化为磁盘 Parquetsave()通过data.to_parquet(...)触发实际计算并落盘load()用dd.read_parquet懒加载。也就是说 Dask 的惰性计算图本身不进入 artifact store进入的是其执行结果。step 边界即计算边界compute_resultstep 中显式调用df.compute()把并行结果收拢为 pandas DataFrame——跨 step 传递时明确在哪里触发计算可以避免下游 step 意外重复计算。用 Numba 加速单机数值代码不是所有瓶颈都需要分布式。对于单机就能跑、但纯 Python 循环太慢的数值代码Numba JIT 是更轻的方案from zenml import step, pipeline import numpy as np from numba import jit import os jit(nopythonTrue) def numba_function(x): return x * x 2 * x - 1 step def load_data() - np.ndarray: return np.arange(1000000) step def apply_numba_function(data: np.ndarray) - np.ndarray: return numba_function(data) pipeline def numba_pipeline(): data load_data() result apply_numba_function(data) # 运行管道 numba_pipeline()jit(nopythonTrue)编译的函数可以像普通 UDF 一样广播到百万级 NumPy 数组上无需任何集群基础设施。分布式引擎集成的重要注意事项教程列出了五条工程约束直接决定方案的稳定性环境准备确保执行环境本地或远端容器已安装 Spark / Ray 等框架。资源管理这些框架会自己管理线程池/集群资源需要与 ZenML 的编排资源分配协调避免超卖。错误处理尤其要写好 Spark session 停止、Ray runtime 关闭的清理逻辑防止容器退出时残留进程。数据 I/O思考数据如何进出分布式 step——对大数据集通常需要借助云存储等中间介质而非通过 artifact store 直接搬运。基础设施规模分布式框架能允许你做分布式计算但底层集群是否扛得住这个计算量仍需自行评估。用 LakeFS 在外部做数据版本化50 GB有时候瓶颈不在算力而在artifact store 本身。哪怕数据没变每次 run 都把 100 GB 的 DataFrame 序列化进 artifact store 依然又慢又贵。对大型、变化缓慢的数据集更好的模式是重型数据留在专用的数据版本化层管道中只传递轻量引用。LakeFS 是这一层的典型实现它在现有对象存储S3、GCS、Azure Blob之上提供类 Git 的分支、提交与回滚。分工很清晰——ZenML 负责管道编排与血缘追踪LakeFS 负责数据。核心模式Pydantic 引用模型定义一个指向 LakeFS 中数据的小 Pydantic 模型让它在 step 之间流转替代数据本身from pydantic import BaseModel class LakeFSRef(BaseModel): repo: str # my-data-repo ref: str # 分支名或 commit SHA path: str # validated/data.parquet endpoint: str # https://lakefs.example.com每个 step 通过 LakeFS 的 S3 兼容 API 直接读写数据ZenML 的 artifact store 只看到这个约 200 字节的引用对象。从源码看这一机制能被自动支撑的原因在于ZenML 内置的 PydanticMaterializer 以BaseModel为关联类型ASSOCIATED_TYPES (BaseModel,)会把模型序列化为 JSON 写入 artifactsave()中调用data.model_dump(modejson)加载时再用model_validate还原。因此LakeFSRef作为 step 输入输出无需任何自定义 materializer。每个 step 的读写示例教程中的核心 stepfrom typing import Annotated from zenml import step step def validate_data(raw_ref: LakeFSRef) - Annotated[LakeFSRef, validated_ref]: # 通过 boto3S3 兼容直接从 LakeFS 读取 df read_from_lakefs(raw_ref) df_clean df.dropna() # 把验证后的数据写回 LakeFS write_to_lakefs(df_clean, raw_ref.repo, raw_ref.ref, validated/data.parquet) # 提交分支——返回不可变 SHA commit_sha commit_branch(raw_ref.repo, raw_ref.ref, Validated data) return LakeFSRef( reporaw_ref.repo, refcommit_sha, # 不可变——保证可复现性 pathvalidated/data.parquet, endpointraw_ref.endpoint, )关键技巧是返回 commit SHA 而非分支名分支名是一个会移动的指针而 commit SHA 是不可变快照。下游 step 无论何时重跑读到的永远是同一份数据——确定性由此免费获得。仓库中的完整可运行示例本仓库的 examples/lakefs_data_versioning 提供了一条端到端落地这条模式的完整管道合成传感器数据生成ingest→ 清洗验证validate→ RandomForest 训练train并附带本地 LakeFS 的 docker-compose 编排。启动方式docker compose up -d # 启动本地 LakeFS端口 8000 pip install -r requirements.txt zenml init python run.py # 运行管道示例的实现细节比教程骨架更完整值得逐一对照引用模型。LakeFSRef 与教程中的模型一致repo/ref/path/endpoint四字段并额外提供一个s3_uri属性拼出s3://{repo}/{ref}/{path}形式的 S3 网关 URI方便任何 S3 兼容工具直接寻址数据。S3 网关读写 SDK 管理的双通道设计。lakefs_utils.py 明确区分了两条通道数据读写全部走 boto3 的 S3 网关get_s3_client()指向LAKEFS_ENDPOINT仓库/分支/提交管理才用 LakeFS SDK。write_parquet_to_lakefs()把 DataFrame 序列化到内存 buffer 再upload_fileobj上传key 为{branch}/{path}read_parquet_from_lakefs()则支持传分支名或 commit SHA作为 ref——这正是不可变引用能生效的前提同一个 object key 在 commit 上下文下内容永不变化。ingest step 的双模式。ingest.py 支持两种运行方式默认生成带约 5% 脏行空值 超温的合成数据写入新分支ingest-{run_id}若传入existing_commit例如 run.py 的--lakefs-commit参数则跳过生成与验证直接返回指向该 commit 下validated/data.parquet的引用——这就是示例 README 中用旧数据快照重新训练能力的实现。validate step 的可复现性闭环。validate.py 在清洗去空值、温度范围过滤后调用commit_branch()拿到 commit SHA 并填入返回的LakeFSRef.ref同时用正则^[0-9a-f]{64}$识别输入 ref——如果传入的已经是 commit SHA复用流程直接透传、跳过验证。验证统计总行数、剔除率等还通过log_metadata写入 ZenML 元数据在 dashboard 中可见。管道组装。training_pipeline.py 把三步串成ingest_data → validate_data → train_model并支持with_options(config_path..., enable_cache...)切换 configs/local.yaml小数据集快速迭代与 configs/remote.yamlK8s / Vertex AI / SageMaker 等远端编排器的 DockerSettings。示例 README 给出的职责划分表可作为判断什么走 artifact store、什么走 LakeFS的通用准则数据存放位置原因传感器数据parquetLakeFS太大不适合 artifact storeLakeFSRef指针ZenML artifact store极小支撑血缘与缓存训练出的模型ZenML artifact store小且受益于版本化指标ZenML metadata可在 dashboard 展示适用边界这个模式在以下条件下收益最大——(1) 数据集大到序列化经 artifact store 成为瓶颈示例中给的经验线是 50 GB 以上(2) 希望在管道版本化之外还要数据层的分支/diff/回滚能力(3) 多条管道或多个团队共享同一批数据集、需要隔离。如何选择正确的扩展策略教程最后给出了五个选型维度可以作为决策清单数据集规模小数据用简单策略起步随数据增长再升级方案。处理复杂度简单聚合可以交给 BigQuery 这类仓库复杂 ML 预处理可能才需要 Spark / Ray。基础设施与资源确认具备支撑分布式处理所需的计算资源。更新频率数据变化越频繁越要考虑重处理成本这正是重数据外置 引用传递模式的价值所在。团队技术栈选择团队熟悉或能快速上手的框架。总原则是从简单开始按需升级。ZenML 管道中 step 与 step 之间的数据流是松耦合的——同一张表今天可以从本地 Parquet 读明天换读 BigQuery后天指向一个 LakeFS commit只要输入输出的类型契约不变管道其余部分无需改动。通过这一系列渐进式策略你可以把 ZenML 管道扩展到任意规模的数据集让 ML 工作流在数据增长时依然保持高效与可控。若要深入自定义 Dataset 类与更复杂的数据流管理可继续参考教程中的自定义 Dataset 类章节。【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进