PySpark チートシート: OVHcloud Data Platform 上の基本的および高度な関数
このチートシートは、OVHcloud Data Platform 上のデータ処理のための基本的および高度な PySpark 関数の簡潔なガイドを提供します
目的
このチートシートは、OVHcloud Data Platform 上のデータ処理のための基本的および高度な PySpark 関数の簡潔なガイドを提供します。これは、2025年1月のNYCイエローキャブ旅行記録(yellow_tripdata_2025_01.parquet、約3.5Mレコード)とタクシーゾーンルックアップテーブル(taxi_zone_lookup.csv、265レコード)を実践的な例として使用します。
中級ユーザーを対象としており、複雑な変換と分析のためのコア関数(filter、select、groupBy、join、udf)と高度な関数(window、pivot、approx_count_distinct、collect_list、explode、regexp_replace)をカバーしています。これらの例は実用的な応用を示しており、このチートシートはどのようなデータセットにも役立つ多才なリファレンスになります。
必要条件
始める前に、以下を確認してください:
- データセット:(
yellow_tripdata_2025_01.parquet、taxi_zone_lookup.csv)が、Connectors に利用可能で、Lakehouse Manager からアクセス可能であること。これらのデータセットは、公式のNYC TLC Trip Record Data ウェブサイトからダウンロードできます。 - ノートブック:PySpark を有効にした Jupyter ノートブック。
セットアップ手順
- Connectors:「NYC-taxi」というソースを作成し、ファイル(
yellow_tripdata_2025_01.parquet、taxi_zone_lookup.csv)をアップロードし、そのスキーマを抽出します。 - Lakehouse Manager:Lakehouse Manager に対応するテーブルを作成します。
- DPE(Data Processing Engine):これらのテーブル内にデータが読み込まれていることを確認します。
- ノートブック:OVHcloud Data Platform 内で PySpark Jupyter ノートブックを開始します。
PySpark 関数 チートシート
ステップ 1: Spark の初期化とデータの読み込み
コードブロック
このコードブロックで使用される関数:
-
SparkSession.builder.appName().getOrCreate():- 何: PySpark セッションと DataFrame API を初期化します。これは Spark 機能を使用するためのエントリーポイントです。
- 利点: 分散データ処理のための環境を設定します。
-
cn_prim = connect("dwh/default_dataset/"):- 何:
forepaas.dwhライブラリを使用して、指定された Lakehouse データセットに接続します。 - 利点: Lakehouse に保存されたテーブルをクエリや操作することを可能にします。
- 何:
-
cn_prim.query(sql_query):- 何:
cn_prim接続を介してアクセス可能なテーブルに対して SQL クエリを実行します。データを Spark DataFrame として取得します。 - 利点: Connectors と Lakehouse から Spark DataFrame にデータを読み込むための簡単な方法を提供します。
- 何:
-
DataFrame.cache():- 何: DataFrame を最初に計算したときにメモリにキャッシュするようにマークします。この DataFrame に対する後のアクションはキャッシュから読み込みます。
- 利点: 同じ DataFrame に対する操作を大幅に高速化します。特に、大規模なデータセット(約 3.5M レコード)に対して反復アルゴリズムや複数の変換を行う場合に有用です。
-
DataFrame.count():- 何: 計算をトリガーし、DataFrame の行の総数を返します。
- 利点: データが読み込まれたことを迅速に確認し、レコード数を確認するために使用されます。
このコードブロックの出力:
- Spark バージョン(例:
Spark Version: 3.4.1) Taxi Records: ~3,500,000(実際のカウントは若干異なる場合があります)Zones Records: 265
ステップ 2: データのクリーンアップと変換
コードブロック
このコードブロックで使用される関数:
-
DataFrame.filter(condition):- 何: 指定された条件に基づいて DataFrame の行をフィルタリングし、条件を満たす行のみを含む新しい DataFrame を返します。
- 利点: データのクリーンアップに不可欠で、無効または関連性のないレコードを削除します。
-
pyspark.sql.functions.col(column_name):- 何: DataFrame の列を参照します。これにより、列に対して様々な変換と操作を適用できます。
- 利点: DataFrame 列を使用した式を構築する方法を提供します。
-
DataFrame.select(columns):- 何: 一連の式(列または列ベースの変換)をプロジェクトし、選択した列のみを含む新しい DataFrame を返します。
- 利点: 列のサブセット化に便利で、関連するフィールドのみを選択することで DataFrame のサイズを削減します。
-
Column.cast(dataType):- 何: 列のデータ型を指定された
dataTypeに変換します。 - 利点: 計算やダウンストリームのプロセスでのデータ型の互換性を確保します(例:
pulocationidをdoubleにキャスト)。
- 何: 列のデータ型を指定された
-
pyspark.sql.functions.regexp_replace(column, pattern, replacement):- 何: 列のテキストデータ内の文字列パターンのすべての出現を指定された置換文字列に置き換えます。
- 利点: データの標準化と文字列フィールドのクリーンアップに優れています(例: "Y" を "Yes" に変換)。
-
DataFrame.withColumn(colName, col):- 何: 指定された式に基づいて新しい列を追加したり、既存の列を置き換えたりして新しい DataFrame を返します。
- 利点: フィーチャーエンジニアリングと新しい派生列のオンデマンド作成に不可欠です。
-
pyspark.sql.functions.unix_timestamp(timestamp_column):- 何: タイムスタンプ文字列またはタイムスタンプ列を Unix タイムスタンプ(1970-01-01 00:00:00 UTC 以降の秒数)に変換します。
- 利点: タイムスタンプの差を計算するなどの数値的な時間計算を容易にします。
-
pyspark.sql.functions.hour(timestamp_column):- 何: タイムスタンプ列から時間成分を抽出します。
- 利点: 時系列分析に便利で、データ内の時間帯ごとのパターンやトレンドを特定できます。
-
pyspark.sql.functions.when(condition, value).otherwise(other_value):- 何: 条件付きロジックを実装します。
conditionが true の場合、列はvalueを取得します。それ以外の場合はother_valueを取得します。チェーン化できます。 - 利点: アウトライアの処理、ビジネスルールの適用、または特定の条件に基づいてデータを分類するのに効果的です(例:
trip_durationのキャップ)。
- 何: 条件付きロジックを実装します。
このコードブロックの出力:
Cleaned Records: ~2,700,000 – ~2,800,000(データ品質に基づいて実際のカウントは異なる場合があります)。
ステップ 3: UDF(ユーザー定義関数)によるカスタムロジック
コードブロック
このコードブロックで使用される関数:
-
udf(func, returnType):- 何: Python 関数を PySpark のユーザー定義関数(UDF)として登録します。これにより、カスタム Python ロジックを Spark DataFrame 列に適用できます。
- 利点: ネイティブ PySpark 関数に利用できない複雑な計算をカスタマイズできます(例: カスタムロジックを使用して「1分あたりの料金」を計算)。
-
FloatType()(pyspark.sql.typesから):- 何: UDF の戻り値のデータ型を Float(単精度浮動小数点数)として指定します。
- 利点: カスタム関数の出力が Spark DataFrame で正しく型付けされることを確保します。
-
DataFrame.withColumn(colName, col):- 何: (ステップ 2 から再利用)式の結果に基づいて新しい列を追加したり、既存の列を置き換えたりします。
- 利点: ここで、新しく定義された
fare_efficiency_udfを適用してfare_efficiency列を作成します。
このコードブロックの出力:
DataFrame.show(5)出力、transformed_dfの最初の 5 行を表示し、新しく追加されたfare_efficiency列を含みます。- 例:
fare_amountが 15.0 でtrip_durationが 600 秒(10 分)の場合、fare_efficiencyは 1.5 ($/min) となります。
- 例:
ステップ 4: 高度なウィンドウ関数とジョイン
コードブロック
このコードブロックで使用される関数:
-
Window.partitionBy(*cols).orderBy(*cols):- 何: ウィンドウ仕様を定義します。
partitionByは行をグループに分け、orderByは各パーティション内の行の論理的な順序を定義します。 - 利点: ランキング、リード/ラグ、累積合計などの高度な分析操作を有効にするために不可欠です。これらは定義された行のサブセットで動作します。
- 何: ウィンドウ仕様を定義します。
-
pyspark.sql.functions.row_number():- 何: ウィンドウ関数で、ウィンドウ仕様で定義された順序に基づいて、各パーティション内の各行に一意の連続番号を割り当てます。
- 利点: レコードのランキングに最適です(例: メトリクスに基づいてトップ N レコードを特定する、トップファーレート)。
-
DataFrame.join(other_df, on=None, how=None):- 何: 指定されたジョイン条件(
on)とジョインタイプ(how、例: "inner"、"left"、"right")に基づいて 2 つの DataFrame を結合します。 - 利点: 異なるソースからの関連情報を統合することでデータを豊かにします(例: タクシー乗車データをゾーンルックアップデータと結合)。
- 何: 指定されたジョイン条件(
-
DataFrame.withColumnRenamed(existing, new):- 何: 既存の列を名前を変更した新しい DataFrame を返します。
- 利点: スキーマを明確にし、読みやすさを向上させます。特にジョイン後に列名が曖昧になる場合に便利です。
-
DataFrame.drop(*cols):- 何: 指定された列を削除した新しい DataFrame を返します。
- 利点: 不要な列を削除することで DataFrame のサイズと複雑さを管理し、メモリを節約します。
このコードブロックの出力:
Joined Records: ~2,600,000 – ~2,700,000(実際のカウントは異なる場合があります)。DataFrame.show()出力、fare_rankが 3 以下の行を表示し、各区のトップ 3 ファーレートを表示します。
ステップ 5: データの集約とピボット
コードブロック
このコードブロックで使用される関数:
-
DataFrame.groupBy(*cols):- 何: DataFrame を 1 つ以上の指定された列でグループ化し、集約計算の準備をします。
- 利点: 異なるカテゴリや次元に基づいてデータの要約と分析を可能にします。
-
DataFrame.agg(*exprs):- 何: グループ化されたデータに集約関数を適用し、要約統計を計算します。
- 利点: 各グループのカウント、合計、平均などのメトリクスを計算するために使用されます。
-
pyspark.sql.functions.approx_count_distinct(column):- 何: グループ内の一意の項目の近似カウントを返します。HyperLogLog++ アルゴリズムを使用します。
- 利点: 非常に大規模なデータセットでは、
countDistinctよりも大幅に高速でメモリ効率が良く、正確なカウントが必ずしも必要ない場合に便利です。
-
pyspark.sql.functions.count(column):- 何: 列の非 NULL 値の数をカウントするか、
count("*")を使用してグループ内のすべての行をカウントします。 - 利点: 各集約グループのサイズまたは特定の値の発生回数を集計します。
- 何: 列の非 NULL 値の数をカウントするか、
-
DataFrame.pivot(pivot_column):- 何: 指定された列の一意の値を新しい列に変換(ピボット)して DataFrame を回転させます。その後の集約が必要です。
- 利点: 通常はレポートや横断的分析に適した広いテーブルを作成し、カテゴリ間の値を直接比較できます。
-
pyspark.sql.functions.avg(column):- 何: 数値列の平均値を計算します。
- 利点: 各グループ内の定量データの中央値を提供します。
-
DataFrame.orderBy(*cols, ascending=True):- 何: 1 つ以上の列に基づいて DataFrame の行を昇順または降順で並べ替えます。
- 利点: 出力を整理し、データを論理的な順序で提示します。
このコードブロックの出力:
agg_df.show():pickup_borough、unique_zones(近似)およびnum_tripsを表示するテーブル(例: マンハッタンには約 60 の一意のゾーンと約 2M の乗車があります)。pivot_df.show(): 行としてpickup_hourを、列としてpickup_boroughを表示するピボットテーブル、各セルに平均fare_amountが含まれます。
ステップ 6: リストの収集と展開
コードブロック
このコードブロックで使用される関数:
-
pyspark.sql.functions.collect_list(column):- 何: 指定された列内の各グループの非 NULL 値をすべて Python リストに集約する集約関数です。
- 利点: 各要素が元のグループのレコードに対応する配列のような構造を作成するのに便利です。
-
pyspark.sql.functions.explode(array_column):- 何: 配列(リスト)またはマップを含む列を、配列/マップ内の各要素に対応する個別の行に変換します。配列に
N個の要素がある場合、その元の行に対してN個の行を作成します。 - 利点: ネストされたデータ構造をフラット化し、個々の要素を個別のレコードとして処理または表示できます。
- 何: 配列(リスト)またはマップを含む列を、配列/マップ内の各要素に対応する個別の行に変換します。配列に
このコードブロックの出力:
DataFrame.show(10)出力、pickup_boroughとexplodedzone列を含むテーブルを表示し、各区の各一意のゾーンが個別の行を取得します(例: マンハッタンの場合、「マンハッタン | ミッドタウン」、「マンハッタン | アッパーイーストサイド」などの複数の行が表示されます)。
ステップ 7: DataFrame をバケットに保存してテーブルを作成(高度)
コードブロック
このコードブロックで使用される関数(および create_table_from_this_dataframe 内):
-
DataFrame.toPandas():- 何: Spark DataFrame を Pandas DataFrame に変換します。これにより、すべての分散データがドライバーノードに収集されます。
- 利点: Pandas 特有の関数を使用してローカルデータ操作とファイルシリアライゼーション(例:
to_csv)を可能にします。 - 警告: 非常に大規模なデータセットでは、ドライバでメモリ不足エラーが発生する可能性があるため、注意して使用してください。
-
Pandas_DataFrame.to_csv(path_or_buffer, index=False, encoding='utf-8'):- 何: Pandas DataFrame をカンマ区切り値(CSV)ファイルに書き込みます。
- 利点: データを保存と取得に適した標準テキスト形式にシリアライズします。
-
io.BytesIO():- 何: Python の
ioモジュールからのクラスで、メモリ内のバイナリストリームを作成します。これはファイルオブジェクトのように動作します。 - 利点: 物理ファイルとやり取りしているかのようにバイトを書き込みおよび読み取ることができ、ディスクに保存せずにサービスに直接データを転送するのに便利です。
- 何: Python の
-
cn_datastore.create_bucket(bucket_name)(forepaas.dwhから):- 何: OVHcloud Data Platform のデータストア内に新しいストレージバケットを作成します。
- 利点: ファイル、中間処理データ、最終処理データを保存するための専用の場所を提供します。
-
cn_bucket.put(file_path, data, size)(forepaas.dwhから):- 何: 通常は BytesIO バッファから、接続されたバケット内の指定されたパスにデータをアップロードします。
- 利点: 処理済みデータ(例:
pivot_dfからの CSV)を OVHcloud クラウドストレージに永続化します。
-
cn_bucket.list()(forepaas.dwhから):- 何: 接続されたバケット内のファイルとサブディレクトリのリストを取得します。
- 利点: ファイルがバケットに正常にアップロードされたことを確認するために使用されます。
-
dwh_request(path, method, json)(forepaas.dwh.commonから):- 何: OVHcloud Data Platform のバックエンドサービスに直接 HTTP API 呼び出しを行うユーティリティ関数です。
- 利点: データソースの作成、Connectors にファイルの登録、Lakehouse Manager でテーブルのビルドをトリガーするなどのタスクを自動化するための細かい制御を提供します。これらは常に
connectオブジェクトを介して公開されるわけではありません。
このステップのプロセスの説明:
この高度なステップでは、処理済み Spark DataFrame(pivot_df)を OVHcloud Data Platform の Lakehouse に新しいテーブルとして永続化する方法を紹介します。PySpark、Pandas、および直接の API 呼び出しの組み合わせを活用します:
- バケット作成: OVHcloud データストアに新しいバケット(
new_bucket)が既に存在しない場合に作成されます。 - ソース作成: 新しいデータソース(
new_source_bucket)が構成され、このバケットにリンクされ、Connectors がその中のファイルを発見できるようになります。 - データシリアライゼーションとアップロード:
pivot_dfは Pandas DataFrame に変換され、メモリ内バッファ(io.BytesIO)内で CSV 形式にシリアライズされます。この CSV データは新しく作成されたバケットにアップロードされます。 - テーブル登録とビルド: 直接の API 呼び出し(
dwh_request)を介して、アップロードされた CSV ファイルが Connectors で「物理テーブル」として登録され、Lakehouse Manager で「論理テーブル」(taxi_pivot_table)が作成され、登録されたソースからデータがロードされます。
警告: これは、直接の API 相互作用に依存する一時的なソリューションです。テーブルの保存と作成のためのよりネイティブで簡素化された SQL 変換と直接ストレージ機能が将来的に OVHcloud Data Platform SDK に含まれる予定です。これにより、手動の API 呼び出しと Pandas 変換の必要性が減少します。
ステップ 8: テーブルの存在を確認
コードブロック
このコードブロックで使用される関数:
-
cn_prim.query(sql_query):- 何: (ステップ 1 から再利用)確立された
cn_prim接続を介して Lakehouse のテーブルに対して SQL クエリを実行し、結果を Spark DataFrame として取得します。 - 利点: ここで、
taxi_pivot_tableが正常に作成され、Lakehouse でクエリ可能であることを確認するために使用されます。
- 何: (ステップ 1 から再利用)確立された
-
DataFrame.count():- 何: (ステップ 1 から再利用)クエリされた DataFrame の行の総数を返します。
- 利点: 新しく作成されたテーブルに期待されるデータとレコードが含まれていることを確認します。
このコードブロックの出力:
- 成功した場合:
taxi_pivot_table Records: [Number of rows in pivot_df](例: 24 時間が存在する場合はtaxi_pivot_table Records: 24)。 - 失敗した場合: テーブルをクエリできないことを示す
logging.errorからのエラーメッセージ。
ステップ 9: クリーンアップと Spark セッションの停止
コードブロック
このコードブロックで使用される関数:
-
SparkSession.catalog.clearCache():- 何: Spark の内部キャッシュをクリアし、キャッシュされた DataFrame と RDD が占有するメモリを解放します。
- 利点: 複雑な操作の後やキャッシュされたデータが必要なくなった場合にリソースを解放するために不可欠です。長時間実行されるアプリケーションやインタラクティブセッションでメモリ不足を防ぐのに役立ちます。
-
SparkSession.stop():- 何: SparkSession を終了し、SparkContext と関連するすべてのリソースを解放します。
- 利点: Spark アプリケーションのクリーンシャットダウンを確保し、クラスターリソースを解放します。Spark アプリケーションの終了時にこの関数を呼び出すことが重要で、リソースリークを防ぎます。
このコードブロックの出力:
INFO - Spark cache cleared and Spark session stopped.- このメッセージは、キャッシュがクリアされ、Spark セッションが正常に終了したことを確認します。通常、カーネルがアイドル状態であるかシャットダウンしたことを示すノートブック環境からのさらに多くの出力を確認できます。
結論と次のステップ
このPySparkチートシートでは、OVHcloud Data Platform上でのデータ処理に必要な基本的および高度な関数のハンズオン概要を提供しました。Sparkの初期化方法、データの読み込みとクリーンアップ、UDFを使用したカスタムロジックの適用、ウィンドウ関数とピボットを使用した複雑な集計の実行、そしてLakehouseへのデータ永続化の管理方法を学びました。
このガイドは、PySpark開発を加速するための実用的なリファレンスとして役立ちます。これらのスキルを現実のシナリオに適用し続けるために、以下のアクションをお勧めします。
- チュートリアルを探索する: NYC Taxi Dataset Analysisチュートリアルに深く掘り下げ、エンドツーエンドのデータパイプラインと機械学習モデルを構築してください。
- 知識を深める: すべての関数について詳細な情報を得るために、公式PySparkドキュメントを参照してください。
コーディングとデータ分析を楽しんでください! 🚀
さらに深く学ぶ
トレーニングや技術サポートが必要な場合は、営業担当者にお問い合わせください、またはこのリンクをクリックして見積もりを依頼し、プロフェッショナルサービスの専門家にプロジェクトのカスタム分析を依頼してください。
Data Platformの開発チームと直接質問し、フィードバックを共有し、交流するには、専用のDiscordチャネルをご利用ください。
OVHcloudサービスについてサポートが必要な場合は、ヘルプセンターでリクエストを作成してください。
ユーザーコミュニティに参加してください。

