Lightning Engine を使用する

Lightning Engine は次世代の Apache Spark パフォーマンスであり、パフォーマンス、費用対効果、運用の安定性を大幅に向上させるように設計された独自の機能強化が導入されています。

利点

Lightning Engine のメリットは次のとおりです。

  • データ オペレーションの高速化: メタデータの処理、書き込みワークロード、ベクトル化された I/O など、クラウド ストレージのインタラクションを最適化することで、パフォーマンスを大幅に向上させ、コストを削減します。

  • インテリジェントなクエリ実行: 高度なオプティマイザーの機能強化を活用して、スキャンされるデータを動的に削減し、データ処理を最適化し、より効率的な実行プランを生成して、より高速で費用対効果の高いクエリを実現します。

  • AI ワークロードと ML ワークロードの効率化: GPU ベースのワークロードのクラスタ起動時間を短縮し、AI 向けに最適化されたイメージを使用して安全な環境でのデプロイを簡素化します。

Lightning Engine はパフォーマンスを大幅に向上させますが、具体的な効果はワークロードによって異なります。I/O の制約を受けるオペレーションよりも、Spark Dataframe API、Spark Dataset API、Spark SQL クエリを活用するコンピューティング負荷の高いタスクに最適です。

標準エンジンとの比較

Lightning Engine は、Managed Service for Apache Spark クラスタで Spark ジョブを実行するために使用される標準エンジンに代わるものです。次の表は、Lightning Engine と標準エンジンのアクティベーション プロパティ、ワークロードの適用性、主なメリットを比較したものです。

機能 標準エンジン Lightning Engine
CLI フラグ --engine=default またはフラグを未設定にします。 --engine=lightning
最適な用途 汎用ジョブ、開発、テスト 大幅な高速化を必要とするエンタープライズ規模のワークロード
主なメリット ベースライン パフォーマンス 最適化されたクラウド ストレージの操作、インテリジェントなクエリ実行

要件

Lightning Engine 機能には次の要件が適用されます。

  • イメージ バージョン: Lightning Engine は、Managed Service for Apache Spark イメージ バージョン 2.3.3 以降の 2.3 サブマイナー イメージ バージョン リリースで使用する必要があります。Managed Service for Apache Spark イメージ バージョン 3.0 では Lightning Engine はサポートされていません
  • サポートされているジョブ: Spark、PySpark、SparkSQL、SparkR がサポートされています。標準エンジンは、Lightning Engine クラスタに送信された他のジョブタイプで実行されます。

ネイティブ クエリ実行

ネイティブ クエリ実行(NQE)は、Lightning Engine のオプション コンポーネントであり、特定のジョブに対する高速化をより高いレベルで実現します。これは、Apache GlutenVelox に基づく Google ハードウェア向けに最適化されたネイティブ エンジンで、Spark クエリの一部を JVM の外部で実行することでパフォーマンスを向上させます。

NQE は次のような場合に推奨されます
Spark Dataframe API と Spark Dataset API を活用するコンピューティング負荷の高いタスク、および Parquet、ORC、Apache Iceberg、Delta Lake のファイルとテーブルからデータを読み取る Spark SQL クエリ。出力ファイル形式はパフォーマンスに影響しません。
NQE が推奨されないケース:
Resilient Distributed Datasets(RDD)、ユーザー定義関数(UDF)、ほとんどの Spark ML ライブラリ、ストレージ アクセスによる遅延を伴う I/O バウンド オペレーションに大きく依存するジョブ。

ARM でのネイティブ クエリ実行

Managed Service for Apache Spark は、2.3-ubuntu22-arm イメージを使用して、ARM アーキテクチャの Lightning Engine 内でネイティブ クエリ実行をサポートします。これは、C4A ARM VM の Google Axion プロセッサ用に特別に最適化されています(Managed Service for Apache Spark でサポートされている ARM マシンタイプをご覧ください)。Google Axion C4A インスタンスは、ARM Neoverse V2 コア上に構築された ARM ベースの VM であり、費用対効果を大幅に向上させます。

Google Axion ARM インスタンスで NQE を実行すると、ハードウェア レベルの価格性能とコアあたりの高いメモリ帯域幅がソフトウェア レベルのネイティブ ベクトル化クエリ アクセラレーションと組み合わされ、クエリの実行時間と総所有コスト(TCO)が大幅に削減されます。

ARM の主なメリットは次のとおりです。

  • オープンソースの Spark と比較して 2.5 ~ 4 倍の高速化: 分析ワークロードで、オープンソースの Apache Spark と比較して 2.5 ~ 4 倍のパフォーマンス向上を実現します。
  • 費用の最適化: Google Axion C4A インスタンスの費用対効果のメリットを最大限に活用します(同等の x86 インスタンスと比較して費用対効果が最大 65% 向上)。
  • オペレーターと形式の完全なパリティ: x86 NQE オペレーターとの完全なパリティ。Cloud Storage Parquet、Apache Iceberg、Delta Lake テーブル形式のネイティブ アクセラレーション。
  • コード変更なし: Spark アプリケーション コードを変更せずに、標準の構成フラグを使用して ARM クラスタでシームレスに有効化できます。

要件

ネイティブ クエリ実行機能には、次の要件が適用されます。

  • 実行エンジン: NQE は、クラスタの作成時に Lightning エンジンが有効になっているクラスタでのみ使用できます。

  • サポートされているジョブ: Spark、PySpark、SparkSQL、SparkR がサポートされています。標準エンジンは、Lightning Engine クラスタに送信された他のジョブタイプで(NQE なしで)実行されます。

  • GPU とアクセラレータなし: GPU アクセラレータで送信された NQE 対応ジョブは失敗します(ただし、NQE なしで Lightning Engine を利用できます)。

  • データ型: 次のデータ型の入力はサポートされていません。

    • バイト: ORC と Parquet
    • 構造体、配列、マップ: Parquet

アーキテクチャとシステム要件

NQE の要件は、プロセッサ アーキテクチャによって異なります。

アーキテクチャ サポートされているマシンタイプ サポートされているオペレーティング システムとイメージ
x86 Intel と AMD のマシン ファミリー Debian-12Ubuntu-22(イメージ バージョン 2.3.3 以降の 2.3 リリース)
ARM C4A マシンシリーズ(Google Axion ARM VM) 2.3-ubuntu22-arm イメージ(バージョン 2.3.3 以降の 2.3 リリース)

料金

料金については、Managed Service for Apache Spark の料金をご覧ください。

Lightning Engine クラスタを作成する

このセクションでは、クラスタに送信された Spark ジョブで Lightning Engine を有効にする Managed Service for Apache Spark クラスタを作成する方法について説明します。

クラスタの作成時にクラスタでネイティブ クエリ実行(NQE)を有効にすることも、クラスタに送信された特定の Spark ジョブに対して後で NQE を有効にすることもできます。

始める前に

  1. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  2. Verify that you have the permissions required to complete this guide.

  3. Verify that billing is enabled for your Google Cloud project.

  4. Enable the Dataproc API.

    Roles required to enable APIs

    To enable APIs, you need the serviceusage.services.enable permission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.

    Enable the API

  5. Google Cloud CLI をインストールします。

  6. フェデレーション ID(連携 ID)を使用するように gcloud CLI を構成します。

    詳細については、連携 ID を使用して gcloud CLI にログインするをご覧ください。

  7. gcloud CLI を初期化するには、次のコマンドを実行します。

    gcloud init

必要なロール

Managed Service for Apache Spark クラスタを作成してクラスタにジョブを送信するには、特定の Identity and Access Management(IAM)ロールが必要です。組織のポリシーによっては、クラウド プロジェクトのオーナーまたはサービス管理者が、これらのロールをユーザーまたはサービス アカウントにすでに付与している場合があります。ロールの付与を確認するには、ロールを付与する必要がありますか?をご覧ください。

ロールの付与の詳細については、プロジェクト、フォルダ、組織へのアクセス権の管理をご覧ください。

ユーザーロール

Managed Service for Apache Spark クラスタの作成に必要な権限を取得するには、次の IAM ロールを付与するよう管理者に依頼してください。

サービス アカウントのロール

Compute Engine のデフォルト サービス アカウントに Managed Service for Apache Spark クラスタを作成するために必要な権限を付与するには、プロジェクトに対する Dataproc ワーカー roles/dataproc.worker)IAM ロールを Compute Engine のデフォルト サービス アカウントに付与するよう管理者に依頼してください。

クラスタを作成する

次の例では、 Google Cloud コンソール、Google Cloud CLI、Dataproc API、Python 用 Cloud クライアント ライブラリ、または Terraform を使用して Lightning Engine クラスタを作成する方法を示します。GoJavaNode.js の Cloud クライアント ライブラリを使用して、Lightning Engine を有効にしたクラスタを作成することもできます。

コンソール

  1. [クラスタの作成] ページを開きます。
  2. [その他の構成] をクリックして、セクションを開きます。
  3. [カスタマイズとその他] を編集します。
  4. 表示されたパネルで、[Lightning Engine を有効にする] チェックボックスがオンになっていることを確認します。
  5. 省略可: Spark ジョブでネイティブ実行ランタイムをデフォルトで有効にするには、[ネイティブ実行を有効にする] チェックボックスをオンにします。
  6. [保存] をクリックします。
  7. 必要に応じて、他のクラスタ設定を構成します。
  8. [クラスタを作成] をクリックします。

ARM クラスタを作成する:

Google Axion とネイティブ クエリ実行を使用して ARM クラスタを作成するには:

  1. [ワーカーの構成] で、C4A マシンシリーズを選択し、ARM マシンタイプを選択します。
  2. [イメージ] で、2.3-ubuntu22-arm イメージを選択します。
  3. [Additional configuration] > [Customization & Other] で、[Enable Lightning Engine] と [Enable Native Execution] が選択されていることを確認します。
  4. [クラスタを作成] をクリックします。

gcloud

  1. Lightning Engine を有効にしてクラスタを作成するには、--engine=lightning フラグを指定して gcloud dataproc clusters create コマンドを実行します。詳細については、gcloud CLI を使用してクラスタを作成するをご覧ください。

    gcloud dataproc clusters create CLUSTER_NAME \
        --region=REGION \
        --engine=lightning \
        --image-version=2.3
    
  2. 省略可: Spark ジョブでネイティブ実行ランタイムをデフォルトで有効にするには、spark:spark.dataproc.lightningEngine.runtime=native プロパティを含めます。

    gcloud dataproc clusters create CLUSTER_NAME \
        --region=REGION \
        --engine=lightning \
        --image-version=2.3 \
        --properties='spark:spark.dataproc.lightningEngine.runtime=native'
    
  3. NQE を使用して ARM クラスタを作成する: Google Axion C4A ARM ノードと NQE が有効になっている Managed Service for Apache Spark クラスタを作成するには:

    gcloud dataproc clusters create CLUSTER_NAME \
        --region=REGION \
        --image-version=2.3-ubuntu22-arm \
        --master-machine-type=c4a-standard-16 \
        --worker-machine-type=c4a-standard-16 \
        --engine=lightning \
        --properties='spark:spark.dataproc.lightningEngine.runtime=native'
    

API

Lightning Engine を有効にしてクラスタを作成するには、clusters.create リクエストを送信します。詳細については、REST API を使用してクラスタを作成するをご覧ください。

  1. リクエストの本文で、engine フィールドを LIGHTNING に設定します。

    {
      "projectId": "PROJECT_ID",
      "clusterName": "CLUSTER_NAME",
      "config": {
        "engine": "LIGHTNING",
        "gceClusterConfig": {},
        "softwareConfig": {
          "imageVersion": "2.3"
        }
      }
    }
    
  2. 省略可: すべてのジョブでネイティブ実行ランタイムをデフォルトで有効にするには、spark:spark.dataproc.lightningEngine.runtime プロパティを含めます。

    {
      "projectId": "PROJECT_ID",
      "clusterName": "CLUSTER_NAME",
      "config": {
        "engine": "LIGHTNING",
        "gceClusterConfig": {},
        "softwareConfig": {
          "imageVersion": "2.3",
          "properties": {
            "spark:spark.dataproc.lightningEngine.runtime": "native"
          }
        }
      }
    }
    

Python

  1. Lightning Engine を有効にしてクラスタを作成するには、create_cluster メソッドを使用して、クラスタ構成の engine フィールドを LIGHTNING に設定します。詳細については、Python でクラスタを作成するをご覧ください。

    from google.cloud import dataproc_v1
    
    def create_lightning_cluster(project_id, region, cluster_name):
        client_options = {"api_endpoint": f"{region}-dataproc.googleapis.com:443"}
        cluster_client = dataproc_v1.ClusterControllerClient(client_options=client_options)
    
        cluster = {
            "project_id": project_id,
            "cluster_name": cluster_name,
            "config": {
                "engine": "LIGHTNING",
                "software_config": {
                    "image_version": "2.3-debian12",
                },
            }
        }
    
        operation = cluster_client.create_cluster(
            project_id=project_id,
            region=region,
            cluster=cluster
        )
        result = operation.result()
        print(f"Cluster created successfully: {result.cluster_name}")
    
  2. 省略可: Spark ジョブでネイティブ実行ランタイムをデフォルトで有効にするには、spark:spark.dataproc.lightningEngine.runtime プロパティを含めます。

    from google.cloud import dataproc_v1
    
    def create_lightning_native_cluster(project_id, region, cluster_name):
        client_options = {"api_endpoint": f"{region}-dataproc.googleapis.com:443"}
        cluster_client = dataproc_v1.ClusterControllerClient(client_options=client_options)
    
        cluster = {
            "project_id": project_id,
            "cluster_name": cluster_name,
            "config": {
                "engine": "LIGHTNING",
                "software_config": {
                    "image_version": "2.3-debian12",
                    "properties": {
                        "spark:spark.dataproc.lightningEngine.runtime": "native"
                    }
                }
            }
        }
    
        operation = cluster_client.create_cluster(
            project_id=project_id,
            region=region,
            cluster=cluster
        )
        result = operation.result()
        print(f"Cluster created successfully: {result.cluster_name}")
    

Terraform

  1. google_dataproc_cluster リソース構成で、engine 引数を LIGHTNING に設定します。
  2. 詳細と詳細オプションについては、google_dataproc_cluster リソースの公式 Terraform ドキュメントをご覧ください。

クラスタ エンジンを確認する

コンソール

  1. Google Cloud コンソールで、[クラスタの詳細] ページに移動します。
  2. Lightning Engine の値が [エンジン] フィールドに表示されていることを確認します。
  3. ネイティブ クエリ実行を有効にした場合は、[ネイティブ実行] フィールドに native が表示されていることを確認します。

gcloud

  1. エンジンと NQE(有効になっている場合)を確認するには、gcloud dataproc clusters describe コマンドを実行します。

    gcloud dataproc clusters describe CLUSTER_NAME --project=PROJECT_ID --region=REGION
    
  2. 出力で engine プロパティと lightningEngine.runtime プロパティを確認します。

    clusterName: lightning-engine-cluster
    engine: lightningEngine
    lightningEngine.runtime: native
    

Lightning Engine を使用してジョブを送信する

クラスタの作成時に Lightning Engine を有効にした場合、クラスタに Spark ジョブを送信すると、ジョブで Lightning Engine がデフォルトで有効になります。

ジョブのネイティブ クエリ実行を有効にする

Lightning Engine クラスタの作成時にネイティブ クエリ実行(NQE)を有効にした場合、特定のジョブで NQE を無効にする場合を除き、すべての Spark ジョブは NQE が有効な状態で実行されます。

Lightning Engine クラスタの作成時に NQE を有効にしなかった場合は、次の例に示すように、ジョブの送信時にジョブの NQE を有効にできます。

gcloud

Spark ジョブを送信するときにネイティブ クエリ実行を有効にするには、spark.dataproc.lightningEngine.runtime=native プロパティを含めます。

gcloud dataproc jobs submit spark \
    --cluster=CLUSTER_NAME \
    --region=REGION \
    --properties=spark.dataproc.lightningEngine.runtime=native \
    -- ...

API

Spark ジョブを送信するときにネイティブ クエリ実行を有効にするには、リクエストに spark.dataproc.lightningEngine.runtime プロパティを含めます。

{
  "job":{
    "placement":{
      "clusterName": ...
    },
    "sparkJob":{
      "mainClass": ...,
      "properties":{
         "spark.dataproc.lightningEngine.runtime":"native"
      }
    }
  }
}

ジョブのネイティブ クエリ実行を無効にする

Lightning Engine クラスタの作成時にネイティブ クエリ実行(NQE)を有効にした場合、特定のジョブで NQE を無効にしない限り、すべての Spark ジョブが NQE を有効にして実行されます。

次の例に示すように、ジョブを送信するときに、特定の Spark ジョブの NQE を無効にできます。

gcloud

Spark ジョブを送信するときに Lightning Engine クラスタでネイティブ クエリ実行を無効にするには、spark.dataproc.lightningEngine.runtime=default プロパティを含めます。

gcloud dataproc jobs submit spark \
    --cluster=CLUSTER_NAME \
    --region=REGION \
    --properties=spark.dataproc.lightningEngine.runtime=default \
    -- ...

API

Spark ジョブを送信するときに Lightning Engine クラスタでネイティブ クエリ実行を無効にするには、spark.dataproc.lightningEngine.runtime=default プロパティを含めます。

{
  "job":{
    "placement":{
      "clusterName": ...
    },
    "sparkJob":{
      "mainClass": ...,
      "properties":{
         "spark.dataproc.lightningEngine.runtime":"default"
      }
    }
  }
}

ジョブのネイティブ クエリ実行を確認する

Lightning Engine クラスタにジョブを送信したら、ジョブでネイティブ クエリ実行が有効になっていることを確認できます。

コンソール

  1. Google Cloud コンソールで、[ジョブ] ページに移動します。
  2. ジョブ ID をクリックして、[ジョブの詳細] ページを開きます。
  3. [ネイティブ実行] フィールドに native が表示されていることを確認します。

gcloud

  1. gcloud dataproc jobs describe コマンドを実行します。

    gcloud dataproc jobs describe JOB_ID --project=PROJECT_ID --region=REGION
    
  2. 出力の [プロパティ] セクションで lightningEngine.runtime を確認します。

    lightningEngine.runtime: native
    

構成パラメータ

次の表に、Lightning Engine とネイティブ クエリ実行の主な構成パラメータを示します。

パラメータ名 説明 該当するエンジン デフォルト値 デフォルト値(Lightning Engine) ユーザーによるオーバーライド可能(ジョブレベル) 範囲
--engine クラスタの作成時にエンジンを選択するクラスタレベルの設定。 クラスタ全体 default lightning × クラスタ
spark:spark.dataproc.lightningEngine.runtime クラスタの作成時に Lightning エンジン ランタイムを選択するクラスタレベルの設定。 Lightning のみ default default × クラスタ
spark.dataproc.lightningEngine.runtime Lightning Engine 内でネイティブ クエリ実行(NQE)を有効または無効にします。 Lightning のみ default default はい。native または default に設定できます。 ジョブ

制限事項

次のシナリオでネイティブ クエリ実行を有効にすると、例外、Spark の非互換性、ワークロードのデフォルトの Spark エンジンへのフォールバックが発生する可能性があります。

フォールバック

次のシナリオでネイティブ クエリを実行すると、ワークロードが Spark 実行エンジンにフォールバックする可能性があります。

  • ANSI: ANSI モードが有効になっている場合、実行は Spark にフォールバックします。
  • 大文字と小文字を区別するモード: ネイティブ クエリ実行では、Spark のデフォルトの大文字と小文字を区別しないモードのみがサポートされます。大文字と小文字を区別するモードが有効になっている場合、正しくない結果が生じる可能性があります。
  • パーティション分割テーブル スキャン: ネイティブ クエリ実行は、パスにパーティション情報が含まれている場合にのみ、パーティション分割テーブル スキャンをサポートします。それ以外の場合、ワークロードは Spark 実行エンジンにフォールバックします。

互換性のない動作

ネイティブ クエリ実行を次のケースで使用すると、互換性のない動作や誤った結果が発生する可能性があります。

  • JSON 関数: ネイティブ クエリ実行では、一重引用符ではなく二重引用符で囲まれた文字列がサポートされます。単一引用符を使用すると、誤った結果が返されます。get_json_object 関数を含むパスで * を使用すると、NULL が返されます。
  • Parquet 読み取り構成:
    • ネイティブ クエリの実行では、spark.files.ignoreCorruptFilestrue に設定されている場合でも、デフォルトの false 値に設定されているものとして扱われます。
    • ネイティブ クエリ実行は spark.sql.parquet.datetimeRebaseModeInRead を無視し、Parquet ファイルの内容のみを返します。従来のハイブリッド カレンダーと先発グレゴリオ暦の違いは考慮されません。Spark の結果は異なる場合があります。
  • NaN: 対象外です。たとえば、数値比較で NaN を使用すると、予期しない結果が生じる可能性があります。
  • Spark カラム型読み取り: Spark カラム型ベクトルがネイティブ クエリ実行と互換性がないため、致命的なエラーが発生する可能性があります。
  • スピル: シャッフル パーティションを大きな数に設定すると、ディスクへのスピル機能が OutOfMemoryException をトリガーする可能性があります。この場合は、パーティションの数を減らすことで、この例外を解消できます。

次のステップ