注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。
軽量トランスフォームは、スキーマのないデータセットのファイルを処理できます。スキーマのないデータセットのファイルを処理するには、my_input.filesystem().ls() でファイルを一覧表示します。
.filesystem().ls() はスキーマのないデータセットで利用可能ですが、.path()、.pandas()、.polars()、.arrow()、および .filesystem().files() はスキーマのあるデータセットでのみ利用可能です。
次のコードは、スキーマのないデータセットのファイルを処理する軽量トランスフォームの例です。
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23from transforms.api import incremental, Input, Output, transform import polars as pl from concurrent.futures import ThreadPoolExecutor @incremental() @transform.using(my_input=Input("my-input"), my_output=Output('my-output')) def my_incremental_transform(my_input, my_output): fs = my_input.filesystem() files = list(fs.ls(glob="*.csv")) def process_file(dataset_file): file_path = dataset_file.path # ファイルにアクセスします with fs.open(file_path, "rb") as f: # <do something with the file> # データをデータフレームとして返します with ThreadPoolExecutor() as executor: polars_dataframes = list(executor.map(process_file, files)) # すべての DF を和集合として1つにまとめます combined_df = pl.concat(polars_dataframes) out.write_table(combined_df)
次の例は、Excel ファイルをパースする方法を示しています。
Copied!1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46from transforms.api import transform, Input, Output import tempfile import shutil import polars as pl import pandas as pd from concurrent.futures import ThreadPoolExecutor @transform.spark.using( my_output=Output("/path/tabular_output_dataset"), my_input=Input("/path/input_dataset_without_schema"), ) def compute(my_input, my_output): # 各ファイルをパースします # 指定されたファイルシステムを使って、指定されたパスの Excel ファイルを開きます def read_excel_to_polars(fs, file_path): with fs.open(file_path, "rb") as f: with tempfile.TemporaryFile() as tmp: # ソースデータセットからローカルファイルシステムにファイルをコピーして貼り付けます shutil.copyfileobj(f, tmp) tmp.flush() # shutil.copyfileobj はフラッシュしません # Excel ファイルを読み込みます(この時点でファイルはシーク可能です) pandas_df = pd.read_excel(tmp) # Integer 型の列がある場合は、文字列型の列に変換します pandas_df = pandas_df.astype(str) # pandas データフレームを polars データフレームに変換します return pl.from_pandas(pandas_df) fs = my_input.filesystem() # 入力データセット内のすべてのファイルを一覧表示します files = [f.path for f in fs.ls()] def process_file(curr_file_as_row): # print(curr_file_as_row) return read_excel_to_polars(fs, curr_file_as_row) def union_polars_dataframes(dfs): return pl.concat(dfs) # すべての DF を和集合として1つにまとめます with ThreadPoolExecutor() as executor: polars_dataframes = list(executor.map(process_file, files)) combined_df = union_polars_dataframes(polars_dataframes) my_output.write_table(combined_df)