with_fs_mount デコレーターを使用して、Alibaba Cloud OSS を MaxFrame 内の分散ストレージとしてマウントし、利用します。FS Mount は、大規模データ処理向けに外部データへの安定したファイルシステムレベルのアクセスを提供します。
ユースケース
FS Mount は、OSS などの永続的なオブジェクトストレージと連携する MaxFrame ジョブにおけるビッグデータ分析に最適です。たとえば、以下のような用途があります。
-
OSS から生データをロード・クリーニング・処理する。
-
中間結果を OSS に書き込み、ダウンストリームのタスクが利用できるようにする。
-
トレーニング済みモデルファイルや設定ファイルなどの静的リソースを共有する。
従来の pd.read_csv("oss://...") のような読み取り/書き込み方法は、分散環境において SDK のパフォーマンスやネットワークオーバーヘッドに制限されます。一方、FS Mount を使用すると、OSS 上のファイルをローカルディスク上のファイルのようにアクセスでき、開発効率が大幅に向上します。
操作手順
サービスを有効化して権限を付与する
-
OSS を有効化し、バケットを作成します。
-
OSS コンソールにログインします。
-
左側のナビゲーションウィンドウで、バケット をクリックします。
-
バケット ページで、バケットの作成 をクリックします。
この例では、バケット名は
xxx-oss-test-shです。
-
-
MaxCompute 用の RAM ロールを作成し、実行環境へのアクセス権限を付与します。
-
RAM コンソールにログインします。
-
左側のナビゲーションバーで、 を選択します。
-
ロール ページで、ロールの作成 をクリックします。
-
ロールの作成 ページの右上隅で、サービスリンクロールを作成 をクリックします。
-
ロールの作成 ページで、信頼プリンシパルタイプ を クラウドサービス に設定します。
-
[Principal Type] で、MaxCompute を選択します。
-
権限管理 タブで、新規権限付与 をクリックします。権限付与の追加 パネルで、ロールに付与するポリシーを選択し、権限を付与 をクリックします。
以下のポリシーを選択します。
-
AliyunOSSFullAccess:OSS を管理する権限を付与します。
-
AliyunMaxComputeFullAccess:MaxCompute を管理する権限を付与します。
-
-
-
with_fs_mount を使用した OSS のマウント
-
推奨:ロール ARN による認証
from maxframe.udf import with_fs_mount @with_fs_mount( "oss://oss-cn-xxxx-internal.aliyuncs.com/xxx-oss-test-sh/test/", "/mnt/oss_data", storage_options={ "role_arn": "acs:ram::xxx:role/maxframe-oss" }, ) def _process(batch_df): import os if os.path.exists('/mnt/oss_data'): print(f"Mounted files: {os.listdir('/mnt/oss_data')}") else: print("/mnt/oss_data not mounted!") return batch_df * 2 -
非推奨:認証情報のハードコード
この方法はテスト目的でのみ使用し、本番環境では推奨されません。
storage_options={ "access_key_id": "LTAI5t...", "access_key_secret": "Wp9H..." }重要AccessKey をハードコードしないでください。
role_arnを使用することで、システムが自動的に一時的な STS トークンをリクエストし、AccessKey ペアの漏洩を防止できます。
with_running_options を使用したリソース割り当ての制御
with_running_options デコレーターを使用して、タスクに割り当てる CPU およびメモリリソースを指定します。
from maxframe.udf import with_running_options
@with_running_options(engine="dpe", cpu=2, memory=16)
@with_fs_mount(...)
def _process(batch_df):
...
|
パラメーター |
推奨値 |
説明 |
|
|
固定 |
FS Mount は現在、DPE エンジンのみをサポートしています。 |
|
|
1~4 |
I/O 負荷の高いタスクや展開処理が多いタスクでは、この値を増やしてください。 |
|
|
8 GB から開始 |
大容量ファイルをロードする場合は、16 GB 以上を推奨します。 |
使用例
推奨パターン:データをバッチ単位で処理します。
大規模データ処理のシナリオでは、MaxFrame の apply_chunk 機能を使用して、入力データをバッチ単位で処理します。
MaxFrame セッションの作成
import os
from odps import ODPS
from maxframe import new_session
from maxframe.udf import with_fs_mount, with_running_options
# ODPS クライアントを初期化します。
# AccessKey ID および AccessKey Secret の文字列をハードコードする代わりに、
# ALIBABA_CLOUD_ACCESS_KEY_ID および ALIBABA_CLOUD_ACCESS_KEY_SECRET
# 環境変数を設定することを推奨します。
o = ODPS(
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
project='<your-project>',
endpoint='https://service.cn-<region>.maxcompute.aliyun.com/api',
)
# 実行イメージを設定します。
# `maxframe_service_dpe_runtime` イメージには、必要な ossfs2 依存関係が含まれています。
# カスタムイメージを使用する場合は、依存関係をダウンロードし、イメージに含める必要があります。
# パッケージはこのコードブロックの下のリンクから入手できます。
options.sql.settings = { "odps.session.image": "maxframe_service_dpe_runtime"}
# セッションを開始します。
session = new_session(o)
print("LogView:", session.get_logview_address())
print("Session ID:", session.session_id)
@with_running_options(engine="dpe", cpu=2, memory=8)
@with_fs_mount(
"oss://oss-cn-<region>-internal.aliyuncs.com/wzy-oss-test-sh/test/",
"/mnt/oss_data",
storage_options={
"role_arn": "acs:ram::<uid>:role/maxframe-oss"
},
)
OSSFS 依存パッケージ:ossfs2_2.0.3.1_linux_x86_64.deb
ユーザー定義関数 (UDF) の作成
def _process(batch_df):
import pandas as pd
import os
# ステップ 1:マウントが成功したか確認します。
mount_point = "/mnt/oss_data"
if not os.path.exists(mount_point):
raise RuntimeError("OSS mount failed!")
# ステップ 2:マッピングテーブルや辞書などのデータをロードします。
mapping_file = os.path.join(mount_point, "category_map.csv")
if os.path.isfile(mapping_file):
mapping_df = pd.read_csv(mapping_file)
# ステップ 3:現在のチャンクを処理します。
result = batch_df.copy()
result['F'] = result['A'] * 10
return result
DataFrame の構築と UDF の適用
import maxframe.dataframe as md
data = [[1.0, 2.0, 3.0, 4.0, 5.0], ...]
df = md.DataFrame(data, columns=['A', 'B', 'C', 'D', 'E'])
# apply_chunk を使用して UDF を適用します。
result_df = df.mf.apply_chunk(
_process,
skip_infer=True,
output_type="dataframe",
dtypes=df.dtypes,
index=df.index
)
# 操作を実行し、結果を取得します。
result = result_df.execute().fetch()
skip_infer=True を設定すると型推論がスキップされ、実行速度が向上します。ただし、dtypes および index が正しく渡されていることを確認する必要があります。
トラブルシューティング
マウント状態の確認
_process 関数内にデバッグログを追加します。
import os
print("Mount path exists:", os.path.exists("/mnt/oss_data"))
print("Files in mount:", os.listdir("/mnt/oss_data") if os.path.exists("/mnt/oss_data") else [])
LogView 出力を確認し、次のようなログが出力されていることを確認します。
FS Mount successful! /mnt/oss_data: ['data.csv', 'config.json', 'model.pkl']
Processing batch with shape: (1000, 5)