Ray Data 保存数据 API 完全指南:将 Dataset 写入文件、数据库与数据湖

发布时间:2026/9/19 22:31:16
Ray Data 保存数据 API 完全指南:将 Dataset 写入文件、数据库与数据湖
Ray Data 保存数据 API 完全指南将 Dataset 写入文件、数据库与数据湖【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/rayRay Data 的 Saving Data API 是ray.data.Dataset面向“输出侧”的一整套公开接口覆盖文件格式写出Parquet / CSV / JSON / ORC / TFRecords / NumPy / Images、数据库写入SQL / MongoDB / BigQuery / Snowflake / ClickHouse、数据湖表格式Iceberg / Lance、与其他计算框架的转换Pandas / Dask / Modin / Mars / Daft / Spark以及面向扩展开发者的 Datasink 与 FilenameProvider 抽象。本文以 doc/source/data/api/saving_data.rst 中列出的 API 清单为骨架逐一讲解其签名、参数语义与底层实现帮助你在实际项目中准确选择并正确调用这些写出接口。公共 API 总览Saving Data API 分为“写入”write_*与“转换”to_*两类。写入类将 Dataset 持久化到目标存储转换类将 Dataset 就地转换为其他框架的内存对象或引用列表。下表为文档列出的全部公共 API类别API目标云数仓Dataset.write_bigqueryGoogle BigQuery文件Dataset.write_csvCSV 文件数据库Dataset.write_clickhouseClickHouse转换Dataset.to_daftDaft DataFrame转换Dataset.to_daskDask DataFrame数据湖Dataset.write_icebergApache Iceberg 表文件Dataset.write_images图片文件文件Dataset.write_jsonJSON / JSONL 文件数据湖Dataset.write_lanceLance 数据集转换Dataset.to_marsMars DataFrame转换Dataset.to_modinModin DataFrame数据库Dataset.write_mongoMongoDB文件Dataset.write_numpy.npy文件文件Dataset.write_orcORC 文件转换Dataset.to_pandaspandas DataFrame文件Dataset.write_parquetParquet 文件数据库Dataset.write_sql任意 DB API2 兼容数据库数据库Dataset.write_snowflakeSnowflake转换Dataset.to_sparkSpark DataFrame文件Dataset.write_tfrecordsTFRecord 文件所有write_*方法在 python/ray/data/dataset.py 中定义均以ConsumptionAPI标注属于“消费型”操作执行后触发实际计算并最终统一委托给Dataset.write_datasink完成分布式写出。例如write_parquetL4847构造一个ParquetDatasink后调用self.write_datasink(datasink, ...)其余格式的实现模式与此完全一致。文件格式写出write_parquet / write_csv / write_json / write_orc文件类写出接口共享同一组核心参数理解它们即可举一反三。write_parquet最常用的列式存储写出签名L4847ds.write_parquet( path, *, partition_colsNone, # 按列 Hive 风格分区如 [year, month] filesystemNone, # pyarrow 文件系统按 path scheme 自动选择 catalogNone, # 可选 Catalog如 Databricks Unity Catalog try_create_dirTrue, # 自动创建目标目录 arrow_open_stream_argsNone, # 传给 pyarrow open_output_stream 的 kwargs filename_providerNone, # 自定义输出文件名 arrow_parquet_args_fnNone, # 惰性解析的 ParquetWriter 参数工厂 min_rows_per_fileNone, # [实验] 每文件目标最小行数 max_rows_per_fileNone, # [实验] 每文件目标最大行数 ray_remote_argsNone, # 传给 ray.remote 的写任务参数 concurrencyNone, # 最大并发写任务数 num_rows_per_fileNone, # [已弃用] 改用 min_rows_per_file modeSaveMode.APPEND, # overwrite / error / ignore / append **arrow_parquet_args, # 直接传给 ParquetWriter 的选项 )关键语义来自 dataset.py 的参数说明输出文件数量默认等于 Dataset 的 block 数量。要控制文件数先调用ds.repartition(n)调整 block 数再写出。默认文件名{uuid}_{block_idx}.parquet其中uuid是 Dataset 的唯一标识可通过自定义FilenameProvider改变命名规则。文件系统自动选择不传filesystem时依据 path 前缀自动选择 pyarrow 实现例如s3://前缀使用S3FileSystem也可显式传入 pyarrow 文件系统实例以覆盖默认配置。分区写入partition_cols指定后按 Hive 分区风格输出目录结构。arrow_parquet_args_fnvsarrow_parquet_args两者都透传给pyarrow.parquet.ParquetWriter当写参数无法被 pickle如包含闭包或文件句柄或希望按 block 惰性解析时应使用arrow_parquet_args_fn它返回参数字典且优先级高于arrow_parquet_args。特别地partitioning_flavor未指定时 Ray Data 默认使用hiverow_group_size会被传递给ParquetWriter.write_table()。Catalog 集成传入catalog如DatabricksUnityCatalog时path被解释为catalog.schema.table形式的表标识符由 Catalog 解析物理写出位置与凭据此时不允许同时指定filesystem且try_create_dir会直接抛ValueError目录创建需要 bucket 级权限而 vended 凭据通常只有前缀级权限。非原子性警告modeoverwrite会先删除path下的全部数据再写入整个过程不是原子的生产环境需自行保证幂等或使用临时目录。示例import ray ds ray.data.range(100) ds.write_parquet(local:///tmp/data/) # 本地目录 ds.write_parquet(s3://bucket/folder/) # 对象存储按 scheme 自动选择 S3FileSystem ds.write_parquet( /tmp/partitioned/, partition_cols[year], moderay.data.SaveMode.OVERWRITE, )write_json / write_csv面向 pandas 互操作write_jsonL5040写出 JSON 与 JSONL 文件write_csvL5543写出 CSV 文件。两者要求数据集记录可分别转换为 pandas DataFrame / pyarrow Table核心参数与write_parquet一致差异仅在底层写出引擎write_json接受pandas_json_args_fn与**pandas_json_args透传给pandas.DataFrame.to_json系列write_csv接受arrow_csv_args_fn与**begin▁of▁sentence#define arrow_csv_args透传给pyarrow.csv.write_csv。示例来自文档字符串import ray ray.data.range(100).write_json(local:///tmp/data) # 默认 .json 后缀 ray.data.range(100).write_csv(local:///tmp/data) # 默认 .csv 后缀 ray.data.range(100).write_csv(s3://bucket/folder/) # 写 S3write_orc列式 ORC 写出write_orcL5667与write_parquet参数结构基本同构底层透传pyarrow.orc.write_table。需要注意不要把compression放进arrow_open_stream_argsORC 压缩应直接作为write_orc的关键字参数传入见 dataset.py 的说明。其arrow_orc_args_fn/arrow_orc_args与 Parquet 版本语义一致。import ray ray.data.range(100).write_orc(/tmp/data) ray.data.range(100).write_orc(s3://bucket/folder/, compressionzstd)专用格式write_tfrecords / write_numpy / write_images这三个接口针对特定领域的数据形态参数与通用接口有明显差异。write_tfrecords机器学习训练数据write_tfrecordsL5791将每一行写为一条tf.train.Example记录ds.write_tfrecords( path, *, tf_schemaNone, # TensorFlow schema_pb2.Schema定义 Example 结构 filesystemNone, try_create_dirTrue, arrow_open_stream_argsNone, filename_providerNone, min_rows_per_fileNone, ray_remote_argsNone, concurrencyNone, num_rows_per_fileNone, # [已弃用] modeSaveMode.APPEND, )重要限制dataset.py 的警告tf.train.Feature原生只支持 int、float、bytes 三种类型因此该方法只接受包含这些类型的数据集遇到不支持的类型会直接报错。默认输出文件名格式为{uuid}_{block_idx}.tfrecords。write_numpy单列.npy写出write_numpyL6004将 Dataset 的某一列写出为.npy文件column是必填参数import ray ray.data.range(100).write_numpy(local:///tmp/data/, columnid)该方法仅支持可转换为 NumPy 数组的列。除column外其余参数filesystem、filename_provider、min_rows_per_file、ray_remote_args、concurrency、mode等与通用文件接口一致。write_images图片批量写出write_imagesL5465将数据集某一图片列写出为图片文件import ray ds ray.data.read_images(s3://anonymousray-example-data/image-datasets/simple) ds.write_images(local:///tmp/images, columnimage, file_formatpng)column指定包含图片数据的列必填file_format默认png可选格式取决于 Pillow 支持的图片格式见pillow文档的 Image file formats 列表其余参数与通用接口相同默认文件名为{uuid}_{block_idx}.{format}。数据库写出write_sql / write_mongo / write_bigquery / write_snowflake / write_clickhousewrite_sql通用 DB API2 数据库写出write_sqlL6104是数据库类写出的基石面向所有提供 Python DB API2PEP 249兼容连接器的数据库import sqlite3 import ray connection sqlite3.connect(example.db) connection.cursor().execute(CREATE TABLE movie(title, year, score)) dataset ray.data.from_items([ {title: Monty Python and the Holy Grail, year: 1975, score: 8.2}, {title: And Now for Something Completely Different, year: 1971, score: 7.5}, ]) dataset.write_sql( INSERT INTO movie VALUES(?, ?, ?), lambda: sqlite3.connect(example.db), ) result connection.cursor().execute(SELECT * FROM movie ORDER BY year) print(result.fetchall())参数语义sqlINSERT INTO语句参数占位符数量必须与表列数一致connection_factory无参函数每次调用返回一个新的 DB API2 Connection 对象。注意由于写任务分布在多个 Ray worker 上连接必须通过工厂在每个 worker 内创建不能共享跨进程的连接对象底层并行写机制每个 block 的写任务使用 DB API2 的executemany批量插入见 dataset.py。write_snowflake基于 write_sql 的封装write_snowflakeL6172接收table名称与connection_parameters透传给snowflake.connector.connect内部从 Dataset schema 读取列名、自动生成INSERT INTO TBL (col1, col2, ...) VALUES (%s, %s, ...)语句然后复用write_sql完成写入见 dataset.pyconnection_parameters dict( user..., accountABCDEFG-ABC12345, password..., databaseSNOWFLAKE_SAMPLE_DATA, schemaTPCDS_SF100TCL, ) ds ray.data.read_parquet(s3://anonymousray-example-data/iris.parquet) ds.write_snowflake(MY_DATABASE.MY_SCHEMA.IRIS, connection_parameters)write_bigqueryGoogle BigQuery 写出write_bigqueryL6306ds.write_bigquery( project_idmy_project_id, datasetmy_dataset.my_table, # 格式为 dataset_id.table_id max_retry_cnt10, # 单 block 写因限流重试的次数上限 overwrite_tableTrue, # True 覆盖False 追加 )实现细节值得注意dataset.pymax_retry_cnt只针对 BigQuery 自身的限流错误与 Ray 的容错重试无关该方法会强制将写任务的ray_remote_args[max_retries]设为 0以避免 block 级任务重试造成重复写入若用户显式设置了非 0 的max_retries会发出警告dataset参数格式为dataset_id.table_id目标表不存在时会自动创建。write_mongoMongoDB 写出write_mongoL6232基于 pymongoarrow 实现要求数据集可转换为 pyarrow Tableds.write_mongo( urimongodb://username:passwordmongodb0.example.com:27017/?authSourceadmin, databasemy_db, # 必须已存在否则 ValueError collectionmy_collection, # 必须已存在否则 ValueError )限制与行为仅支持 pymongoarrow 类型映射的子集写出不支持的类型会在类型检查阶段失败每条记录作为新文档插入。若记录带_id字段该_id在集合中必须不存在否则写入被拒绝并失败保护既有文档不被修改不带_id时由 MongoDB 自动生成见 dataset.py。write_clickhouseClickHouse 写出write_clickhouseL6380参数最为丰富支持建表语义import pyarrow as pa import ray docs [{title: ClickHouse Datasink test} for key in range(4)] ds ray.data.from_pandas(pd.DataFrame(docs)) user_schema pa.schema([(id, pa.int64()), (title, pa.string())]) ds.write_clickhouse( tabledefault.my_table, # 完全限定表名 dsnclickhousehttp://user:passlocalhost:8123/default, moderay.data.SinkMode.OVERWRITE, # CREATE / APPEND / OVERWRITE schemauser_schema, # 建表/覆盖时必填 table_settingsray.data.ClickHouseTableSettings( engineReplacingMergeTree(), # 默认 MergeTree() order_byid, # 未提供时自动选择 ), )关键参数dataset.pymode使用SinkMode枚举CREATE已存在则失败、APPEND存在则追加不存在则建表、OVERWRITE先删表再重建。CREATE/APPEND在表不存在时必须提供schema或从首个 block 推断OVERWRITE则必须提供schemaschema为 pyarrow.Schema表已存在且仅追加时可省略client_settings/client_kwargs分别对应 ClickHouse 的 server settings 与客户端选项table_settingsClickHouseTableSettings可指定engine默认MergeTree()、order_by覆盖已有表时复用原 ORDER BY否则自动挑选“最佳”列优先时间戳列其次非字符串列最后第一列、partition_by、primary_key、settingsmax_insert_block_rows超大 block 时可设置该值将单次插入拆分为多次小批量插入默认None不拆分DSN 使用 ClickHouse HTTP 连接串格式如clickhousehttp://username:passwordhost:8123/default。数据湖表格式write_iceberg 与 write_lancewrite_iceberg支持 Upsert 与部分覆盖write_icebergL5166将数据写入 Apache Iceberg 表import ray import pandas as pd from ray.data import SaveMode from ray.data.expressions import col docs [{id: i, title: fDoc {i}} for i in range(4)] ds ray.data.from_pandas(pd.DataFrame(docs)) ds.write_iceberg( table_identifierdb_name.table_name, catalog_kwargs{name: default, type: sql}, # 传给 pyiceberg load_catalog() )该接口的亮点dataset.py自动 Schema 演进新列会自动加入表 schemaschema 从写入数据中自动提取三种写模式SaveMode.APPEND默认直接追加不检查重复SaveMode.UPSERT按upsert_kwargs中的join_cols匹配命中则更新整行未命中则插入Ray Data 采用 copy-on-write 策略始终更新匹配键对应的所有列、插入所有新键以换取最优并行度。upsert_kwargs支持join_cols、case_sensitive、branchSaveMode.OVERWRITE替换全部数据或仅替换满足overwrite_filter的行用ray.data.expressions构造如col(date) 2024-10-28、(col(region) US) (col(status) active)overwrite_kwargs透传给 pyiceberg 的table.overwrite()支持case_sensitive、branchsnapshot_properties提交时写入 snapshot 的自定义属性与catalog参数互斥catalog_kwargs与catalog不能同时指定否则抛ValueError。write_lance面向向量检索的 Lance 格式write_lanceL6638写出 Lance 数据集其默认行数边界与前面接口不同ds.write_lance( /tmp/lance_data/, modeSaveMode.OVERWRITE, # CREATE / APPEND / OVERWRITE )schema缺省时从数据推断min_rows_per_file/max_rows_per_file默认分别为 1 Mi 与 64 Mi 行约 1,048,576 / 67,108,864 行data_storage_version默认legacyv1 版本新版本存储效率更高但需要更新的 lance 版本才能读取支持通过table_idnamespace_impl如rest、dirnamespace_properties写入 namespace 管理的表此时path被忽略且仅支持SaveMode.CREATE。转换到其他计算框架to_daft / to_dask / to_mars / to_modin / to_pandas / to_spark文档中的to_*系列方法用于将 Ray Data Dataset 就地转换为其他生态的 DataFrame便于复用既有分析代码Dataset.to_daft转换为 Daft DataFrameDataset.to_dask转换为 Dask DataFrame供 Dask 分布式计算链路消费Dataset.to_mars转换为 Mars DataFrameDataset.to_modin转换为 Modin DataFrameDataset.to_pandas转换为单机 pandas DataFrame数据将物化到本地内存适合结果回传而非大规模计算Dataset.to_spark转换为 Spark DataFrame。这些方法适合“Ray Data 做 ETL / 预处理其他框架做下游分析”的混用场景to_pandas是最常用的一条通常位于处理管线的末端。开发者 API扩展写出能力文档的 Developer APIs 部分面向需要自定义写出逻辑的开发者核心入口是 python/ray/data/datasource 目录中的基类与工具API用途Dataset.to_arrow_refs返回各 block 对应的ObjectRef[pyarrow.Table]列表零拷贝衔接 pyarrow 生态Dataset.to_numpy_refs返回各 block 对应的ObjectRef[np.ndarray]引用列表Dataset.to_pandas_refs返回各 block 对应的ObjectRef[pd.DataFrame]引用列表Datasink自定义数据写出目标的抽象基类所有write_*最终都落到它身上Dataset.write_datasink低层写出入口接收任意Datasink实例配合ray_remote_args、concurrency执行分布式写出datasource.RowBasedFileDatasink面向“逐行生成文件”场景的文件类 Datasink 基类datasource.BlockBasedFileDatasink面向“每个 block 写一个文件”场景的文件类 Datasink 基类datasource.WriteResult写结果类型供 WriteReturnType 使用datasource.WriteReturnType指定 Datasink 写结果类型的枚举datasource.FilenameProvider输出文件命名策略接口filename_provider参数即接收该接口的实现理解 write_datasink 与 SaveMode所有write_*方法最终都调用Dataset.write_datasinkL6764将格式相关的逻辑封装进各自 Datasink。这也是为什么文档将Datasink系列列为“开发者 API”接入新存储格式 实现一个Datasink子类 一行write_datasink调用。写模式的统一枚举定义在 python/ray/data/_internal/savemode.py 的SaveMode中成员值语义CREATEcreate新建数据若已存在则报错APPENDappend追加新数据不修改已有数据默认OVERWRITEoverwrite用新数据替换全部已有数据非原子先删后写IGNOREignore若数据已存在则不写入ERRORerror若数据已存在则抛错UPSERTupsert按主键更新匹配行、插入新行需指定键字段注意个别接口如write_clickhouse使用独立的SinkModeCREATE/APPEND/OVERWRITE见 clickhouse_datasink.pywrite_lance则默认SaveMode.CREATE调用前请核对具体接口文档。源码速查全部公共写出/转换方法定义python/ray/data/dataset.pywrite_parquetL4847、write_jsonL5040、write_icebergL5166、write_imagesL5465、write_csvL5543、write_orcL5667、write_tfrecordsL5791、write_numpyL6004、write_sqlL6104、write_snowflakeL6172、write_mongoL6232、write_bigqueryL6306、write_clickhouseL6380、write_lanceL6638、write_datasinkL6764SaveMode枚举python/ray/data/_internal/savemode.pySinkMode枚举ClickHousepython/ray/data/_internal/datasource/clickhouse_datasink.pyDatasink 基类与文件名策略python/ray/data/datasource 目录各格式 Datasink 实现与测试python/ray/data/tests/datasource如test_parquet.py实践建议绝大多数场景优先选择write_parquet列式存储、生态最广需要与既有 pandas/Dask 生态衔接时使用to_pandas/to_dask写入数仓或数据湖时优先使用各厂商的原生接口write_bigquery/write_snowflake/write_iceberg等它们内置了重试、覆盖与 Schema 管理等语义自定义输出格式或命名规则时从FilenameProvider与Datasink入手即可无需修改 Ray Data 核心代码。【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考