AI Runtimeにデータをロードする
プレビュー
この機能は パブリック プレビュー段階です。
データおよびモデル資産は、大規模言語モデル(LLM)やビジョン・言語モデル(VLM)のディープラーニングおよびポストトレーニング(事後トレーニング)ワークロードにとって不可欠です。AI Runtime を使用すると、すべてのデータおよびモデル資産に Unity Catalog を介してアクセスできます。
- Unity Catalog ボリューム:主に大規模なデータセットや、画像、音声、テキストなどの非構造化ファイルに使用されます。
- Unity Catalog テーブル: 構造化データおよび表形式データに使用され、Spark Connect を介してアクセスされます。
ボリュームとテーブルが Unity Catalog に登録されており、ユーザーまたは Service Principal からアクセスできる必要があります。
非構造化データ用の Unity Catalog ボリューム
Unity Catalog ボリュームは、構造化データ、半構造化データ、非構造化データなど、あらゆる形式の非表形式データへのガバナンスされたアクセスを提供します。AI Runtime では、ボリュームは、大規模データセット、テキスト、モデル資産、およびモデル チェックポイントにアクセスするための主要なメカニズムです。
ユーザーは、ローカル ディスク上のファイルの操作と同様に、なじみのあるファイルシステム操作を使用して、Unity Catalog ボリューム内のファイルのリスト表示、読み取り、および書き込みを行うことができます。
import os
dir_path = "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir"
file_path = os.path.join(dir_path, "test_file")
os.makedirs(dir_path, exist_ok=True)
# Write to the file
with open(file_path, "w") as file:
file.write("Hello, World!")
同様に、シェル操作も同じように機能します。
%sh ls -l /Volumes/<catalog-name>/<schema-name>/<volume-name>
%sh mkdir -p /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir
%sh touch /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/test_file
Unity Catalog のボリュームには、Machine Learning ワークロードに適したいくつかの特徴があります。
- 分散ストレージ : Unity Catalog は分散ストレージによってバックアップされており、AI Runtime のワークロードが、ノートブックと CLI ベースのワークロードの両方から、プラットフォーム全体でデータやモデルアセットを読み書きできるようにします。
- 機械学習 アクセス パターン向けに最適化 : 基盤となるストレージとアクセスパスは、一般的な機械学習ワークロード、特にシーケンシャルな読み取り/書き込みを行う大きなファイル向けに最適化されています。これにより、Unity Catalog はトレーニングデータの読み込み、モデルアセットの読み込み、およびモデルチェックポイントの書き込みに非常に適したツールとなります。
- ファイルシステムのようなアクセス : ユーザーは、ローカルディスク上のファイルの操作と同様に、使い慣れたファイルシステム操作を使用して Unity Catalog ボリューム内のファイルを一覧表示、読み取り、および書き込みできます。
バックグラウンドでの自動 commit により、ユーザーは Unity Catalog ボリュームデータに一貫してアクセスできます。
- 書き込み :AI ランタイムは書き込みを自動的に commit し、同じ Unity Catalog ボリュームにアクセスする他のアプリケーションやワークロードから変更を確認できるようにします。
- 読み取り : AI ランタイムでは、ユーザーによる明示的な更新や同期操作を必要とせず、ボリュームへの変更が自動的に取得されます。
ボリュームのパフォーマンスを調整する
前述のとおり、Unity Catalog ボリュームは分散ストレージを基盤とし、シーケンシャルな読み取りおよび書き込みを行う大規模ファイル向けに最適化されています。
AI ランタイム で最高のパフォーマンスを引き出すためのヒント:
-
データをより大きなファイルに連結する: 可能であれば、データを少数のより大きなファイル(ファイルあたりおおむね 1 GiB ~ 10 GiB)に統合します。これにより、AI Runtime はデータを積極的にプリフェッチし、ほぼ最適なシーケンシャル読み取りパフォーマンスを自動的に達成できます。
-
小さいファイルを扱うワークロードでは、ローカルディスクを使用してください:ワークロードに多数の小さいファイルが含まれる場合は、処理前に並列コピーを使用してファイルをローカルディスクにコピーすることを検討してください。これにより、ボリュームを介して多数の小さいファイルに繰り返しアクセスする際のオーバーヘッドを軽減できます。
sh# Recommended using parallel copy (256 concurrency in this example, you can tune)
#
# This takes only 22 seconds to copy 15,375 150KiB small image files.
%sh cd /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ && find . -type f -print0 | xargs -0 -P 256 -I {} cp --parents "{}" /tmp/
# !!! Avoid doing this !!!
#
# Because the files are copied in serial, this copies the same 15,375 150KiB small image files much more slowly.
# %sh cp -r /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/* /tmp -
Machine Learning ワークロードに
UCVolumeDatasetを使用できます。Unity Catalog ボリュームから効率的なデータアクセスとロードを提供するために、前述の最適化が組み込まれています。詳細については、以下のセクションを参照してください。
非構造化データの読み込み元: UCVolumeDataset
Unity Catalogボリュームに保存されている画像、音声、テキストファイルなどの非構造化データの場合は、databricks.air.dataモジュールのUCVolumeDatasetを使用します。UCVolumeDatasetは、初回アクセス時にボリュームから高速なローカルキャッシュに各ファイルをコピーし、キャッシュされたローカルファイルのパスを生成するPyTorchのIterableDatasetです。手動で実装する必要があるパフォーマンスや分散に関する処理が自動的に処理されます。
- ローカルキャッシュ ファイルは、最初のアクセス時にFUSEマウントからローカルキャッシュディレクトリにコピーされ、その後、キャッシュから提供されるため、マルチエポックトレーニングでボリュームを再度読み込むことはありません。
- 自動パーティショニング
torch.distributedが初期化されると、ファイルはランク全体でパーティション分割され、さらにDataLoaderワーカーに分割されるため、各(rank, worker)ペアは、追加のセットアップなしで重複しないスライスを受け取ります。
UCVolumeDataset および databricks.air.data.DataLoader は databricks-sdk-air パッケージから提供されます。互換性のある torch も同時に取得する data extra を使用してインストールします。
%pip install "databricks-sdk-air[data]"
UCVolumeDataset 未加工のローカルファイルパスを返します。それらのファイルをテンソルにデコードするには、パスのストリームを取り込み、解析ロジックを適用する2番目のIterableDatasetでラップします。これはI/Oと構文解析に関する事項を分離します。
from databricks.air.data import UCVolumeDataset
from torch.utils.data import IterableDataset
from PIL import Image
import torchvision.transforms.functional as TF
class ImageDataset(IterableDataset):
"""Decodes each cached file path from UCVolumeDataset into a tensor."""
def __init__(self, path_dataset: UCVolumeDataset):
self._path_dataset = path_dataset
def __iter__(self):
for local_path in self._path_dataset:
image = Image.open(local_path).convert("RGB")
yield TF.to_tensor(image)
path_dataset = UCVolumeDataset("/Volumes/catalog/schema/my_volume/images")
dataset = ImageDataset(path_dataset)
ラッパーはすでにキャッシュされたローカルパスを受け取るため、解析ステップでボリュームが使用されることはありません。拡張、トークン化、またはフィルタリングのために追加のラッパーをチェインできます。
最適なパフォーマンスを得るには、通常の PyTorch の DataLoader ではなく、UCVolumeDataset と databricks.air.data.DataLoader を組み合わせて使用します。これは AI ランタイム の I/O 向けにチューニングされており、GPU によるコンピュート中にファイルを並列で取得およびキャッシュします。
ボリューム上のチェックポイントモデル
最新のスナップショットからトレーニングを再開したりクラッシュから復旧したりできるようにモデルのチェックポイントを作成するには、ローカルファイルシステムと全く同様に Unity Catalog ボリュームを使用できます。
Databricks では、シングル GPU とマルチ GPU の両方のワークロードでパフォーマンスを向上させるために、分散チェックポイント (DCP) の使用を推奨しています。Databricks エンジニアリングブログの「AI ランタイムにおける高速で耐障害性のある PyTorch トレーニング」を参照してください。
import torch.distributed.checkpoint as dcp
from torch.distributed.checkpoint.state_dict import get_state_dict, set_state_dict
import databricks.air.data
checkpoint_path = "/Volumes/my-catalog/my-schema/my-volume/checkpoints/step_1000"
# Save
model_sd, optim_sd = get_state_dict(model, optimizer)
state_dict = {"model": model_sd, "optim": optim_sd, "step": 1000}
dcp.async_save(
state_dict,
storage_writer=databricks.air.data.UCVolumeWriter(checkpoint_path))
# Load
model_sd, optim_sd = get_state_dict(model, optimizer)
state_dict = {"model": model_sd, "optim": optim_sd}
dcp.load(
state_dict,
storage_reader=databricks.air.data.UCVolumeReader(checkpoint_path))
set_state_dict(
model,
optimizer,
model_state_dict=state_dict["model"],
optim_state_dict=state_dict["optim"],
)
モノリシックな torch.save アプローチも機能します。
-
シングル GPU モデルのチェックポイント作成については、
Python# The monolithic torch.save approach for single GPU chip
# Save
torch.save({"model": model.state_dict(), "opt": optimizer.state_dict()},
"/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt")
# Load
ckpt = torch.load(
"/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt",
weights_only=True)
model.load_state_dict(ckpt["model"])
optimizer.load_state_dict(ckpt["opt"]) -
torchrun を介して起動される分散トレーニングの場合、
Python# The monolithic torch.save approach for multi-GPU distributed training.
# This snippet assumes your launcher has already called
# dist.init_process_group(...).
import os
import torch.distributed as dist
# Save only on rank 0.
if dist.get_rank() == 0:
torch.save({"model": model.state_dict(), "opt": optimizer.state_dict()},
"/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt")
# Wait for rank 0 to finish writing before any rank reads.
dist.barrier()
# Load on ALL ranks (map to current rank's local GPU).
local_rank = int(os.environ["LOCAL_RANK"])
ckpt = torch.load(
"/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt",
map_location=f"cuda:{local_rank}",
weights_only=True)
model.load_state_dict(ckpt["model"])
optimizer.load_state_dict(ckpt["opt"])
表形式のデータを読み込む
Spark Connectを使用して、 Deltaテーブルから表形式の機械学習データをロードします。
シングルノードの場合PySparkメソッドtoPandas()を使用してApache Spark DataFrames Pandas DataFramesに変換し、必要に応じてPySparkメソッドto_numpy()を使用してNumPy形式に変換できます。
Spark Connectは、解析と名前解決を実行時に行うため、コードの動作が変わる可能性があります。Spark ConnectとSpark Classicの比較を参照してください。
Spark Connect は、 Spark SQL 、 Spark上のPandas API 、構造化ストリーミング、 MLlib (DataFrame ベース) など、ほとんどのPySpark APIsサポートします。 最新のサポート対象APIsについては、 PySpark APIリファレンスドキュメントを参照してください。
その他の制限については、 「 サーバレス コンピュートの制限 」を参照してください。
Unity Catalog ボリュームを使用して大規模な Delta テーブルをロードする
toPandas()で変換するには大きすぎる大規模なDeltaテーブルの場合は、データをUnity Catalogボリュームにエクスポートし、 PyTorchまたはHugging Faceを使用して直接ロードします。
# Step 1: Export the Delta table to Parquet files in a UC volume
output_path = "/Volumes/catalog/schema/my_volume/training_data"
spark.table("catalog.schema.my_table").write.mode("overwrite").parquet(output_path)
# Step 2: Load the exported data directly using Hugging Face datasets
from datasets import load_dataset
dataset = load_dataset("parquet", data_files="/Volumes/catalog/schema/my_volume/training_data/*.parquet")
このアプローチは、トレーニング中のSparkのオーバーヘッドを回避し、シングルGPUと分散型トレーニングの両方のワークフローでうまく機能します。