0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

【Managed Service for Apache Spark】Spark MLlibでロジスティック回帰!航空便の遅延予測を自分のプロジェクトで動かす

0
Posted at

Spark の機械学習ライブラリで、フライトが時間どおりに着くかを予測する。クラスタの管理は Google Cloud に任せる。

Managed Service for Apache Spark(旧 Dataproc)を使えば、Spark が動くクラスタをコマンド 1 つで作成できます。当記事では、米国の航空便データを題材に、Spark MLlib のロジスティック回帰で「到着遅延が 15 分未満か」を予測するモデルを作り、保存、復元、評価までを実行します。自分の Google Cloud プロジェクトで実際に動かした結果と、その過程で見つかった注意点をあわせて解説します。

概要

Managed Service for Apache Spark とは

Managed Service for Apache Spark は、Apache Spark と Hadoop のクラスタを作成、実行するためのマネージドサービスです。以前は Dataproc と呼ばれていました。2026年9月現在、コンソールでは「マネージド Apache Spark」と表示されますが、CLI は gcloud dataproc、API は dataproc.googleapis.com のままです。

Spark をオンプレミスで使う場合、Hadoop のインストール、YARN の設定、ノード間の通信設定などを自分で行う必要があります。Managed Service for Apache Spark では、マシンタイプとノード数を指定するだけで、Spark と Jupyter が使えるクラスタが数分で起動します。

当記事で行うこと

ロジスティック回帰は、入力から「ある事象が起きる確率」を求める分類モデルです。当記事では次の設定で使います。

項目 内容
予測したいこと 到着遅延(ARR_DELAY)が 15 分未満か(1:定刻、0:遅延)
入力(特徴量) 出発遅延(DEP_DELAY)、タキシングアウト時間(TAXI_OUT)、飛行距離(DISTANCE)
使うライブラリ Spark MLlib の LogisticRegressionWithLBFGS

当記事は、Python と gcloud コマンドの基本的な操作を前提としています。

構築するもの

gs://spls/gsp271/(ラボ用の公開データ)
        │  gcloud storage cp
        ▼
Cloud Storage  <プロジェクトID>-dsongcp
  ├─ data-science-on-gcp/flights/trainday.csv      学習日/評価日の区分
  └─ data-science-on-gcp/flights/tzcorr/           フライトデータ(CSV)
        │
        ▼
Managed Service for Apache Spark クラスタ ch6cluster
  マスター 1 台 + ワーカー 2 台(e2-standard-4)、Jupyter コンポーネント
        │  PySpark:読み込み → クリーニング → 学習 → 保存 → 評価
        ▼
Cloud Storage  flights/sparkmloutput/model/(保存したモデル)

題材とデータの出典

当記事の題材は、Google Skills のラボ「Machine Learning with Spark on Google Cloud Managed Apache Spark」(ID: GSP271)です。ラボは一時的なプロジェクトで動作しますが、当記事では同じ内容を自分のプロジェクトで実行しました。ラボのコードは、書籍『Data Science on the Google Cloud Platform』(O'Reilly)のリポジトリ GoogleCloudPlatform/data-science-on-gcp の 06_dataproc がもとになっています。

データの元は、米国運輸統計局(Bureau of Transportation Statistics)が公開している米国内線の運航実績です。ラボでは、これにタイムゾーン補正を加えた CSV を使います。

出典:ラボで使用されている Google 提供のサンプルデータ(gs://spls/gsp271/)です。2026年9月の執筆時点では、認証なしで読み取れることを確認しました。一方、書籍側のバケット(gs://data-science-on-gcp/edition2/)は認証なしでは 403 になり、形式も JSON で、ラボのコードとは合いません。ライセンスは明記されていないため、当記事ではファイル自体は再配布せず、取得元のパスと手順のみを記載しています。

オブジェクト 内容
flights/trainday.csv 2015 年の 365 日それぞれが学習日(True)か評価日(False)か。True 263 日、False 102 日
flights/tzcorr/all_flights-00001-of-00025 学習に使うシャード。16,904 行、2015-10-24 と 10-25 の 2 日分
flights/tzcorr/all_flights-00002-of-00025 評価に使うシャード。20,867 行、2015-10-29 と 10-30 の 2 日分
flights/tzcorr/tzcorr.json シャードの列定義(BigQuery のスキーマ形式)

事前準備

環境

  • 課金が有効な Google Cloud プロジェクト
  • Cloud Shell、または gcloud の認証が済んだローカルシェル
gcloud services enable \
  dataproc.googleapis.com \
  compute.googleapis.com \
  storage.googleapis.com

変数の設定

export PROJECT_ID=$(gcloud config get-value project)
export REGION="us-central1"
export BUCKET="${PROJECT_ID}-dsongcp"
export SA_EMAIL="spark-lab@${PROJECT_ID}.iam.gserviceaccount.com"

バケット名は <プロジェクトID>-dsongcp のままにしてください。ラボのコードは BUCKET = PROJECT + '-dsongcp' でバケット名を組み立てるため、この名前にしておけば、コードを変更せずに動きます。

クラスタ用のサービスアカウント

クラスタの VM は、サービスアカウントの権限で Cloud Storage などにアクセスします。ラボでは Compute Engine のデフォルトのサービスアカウント(ラボ環境では編集者ロール付き)が使われますが、当記事では専用のサービスアカウントを作り、必要なロールだけを付与しました。

gcloud iam service-accounts create spark-lab \
  --display-name="GSP271 Spark ML cluster"

gcloud projects add-iam-policy-binding "${PROJECT_ID}" \
  --member="serviceAccount:${SA_EMAIL}" \
  --role="roles/dataproc.worker" \
  --condition=None

roles/dataproc.worker(Dataproc ワーカー)には、クラスタの動作に必要な権限と、Cloud Storage のオブジェクトの読み書き権限が含まれます。今回のデータの読み込みとモデルの保存は、このロールだけで動作しました。

詳細:公式ドキュメント(Managed Service for Apache Spark のサービス アカウント)


データの準備

学習と評価に使うファイルだけをコピーします。

gcloud storage buckets create "gs://${BUCKET}" \
  --location="${REGION}" \
  --uniform-bucket-level-access

SRC=gs://spls/gsp271/data-science-on-gcp/flights
gcloud storage cp "${SRC}/trainday.csv" \
  "gs://${BUCKET}/data-science-on-gcp/flights/"
gcloud storage cp "${SRC}/tzcorr/all_flights-0000[12]-of-00025" "${SRC}/tzcorr/tzcorr.json" \
  "gs://${BUCKET}/data-science-on-gcp/flights/tzcorr/"

ラボと同じ data-science-on-gcp/flights/ の構成でコピーされるため、以降のコードのパスはそのまま使えます。tzcorr/ には全 25 シャード(約 1.8 GB)がありますが、ラボで使うのは 2 シャード(約 11 MB)だけです。

図1 コピー後のバケット。赤枠:学習日の区分 trainday.csv と、フライトデータのフォルダ tzcorr/
SS-01.png


クラスタを作成する

ラボの変更点を反映したコマンド

ラボでは、リポジトリの create_cluster.sh を nano で編集してから実行します。編集内容は次の 3 点です。

ラボでの編集 理由
--zone ${REGION}-a を削除 ゾーンを固定せず、自動で選ばせる
マシンタイプを n1-standard-4 から e2-standard-4 に変更 新しい汎用マシンタイプに合わせる
--public-ip-address を追加 VM に外部 IP を付ける(後述)

当記事では、編集後のスクリプトと同じ内容を gcloud で直接実行しました。専用のサービスアカウントを使うため、--service-account を追加しています。

git clone --depth 1 https://github.com/GoogleCloudPlatform/data-science-on-gcp/
sed "s/CHANGE_TO_USER_NAME/dataproc/g" data-science-on-gcp/06_dataproc/install_on_cluster.sh \
  > install_on_cluster.sh
gcloud storage cp install_on_cluster.sh "gs://${BUCKET}/flights/dataproc/install_on_cluster.sh"

gcloud dataproc clusters create ch6cluster \
  --region="${REGION}" \
  --enable-component-gateway \
  --master-machine-type=e2-standard-4 --master-boot-disk-size=500 \
  --num-workers=2 \
  --worker-machine-type=e2-standard-4 --worker-boot-disk-size=500 \
  --optional-components=JUPYTER \
  --initialization-actions="gs://${BUCKET}/flights/dataproc/install_on_cluster.sh" \
  --scopes=https://www.googleapis.com/auth/cloud-platform \
  --public-ip-address \
  --service-account="${SA_EMAIL}"

install_on_cluster.sh は初期化アクションです。全ノードで Python パッケージを入れ、マスターノードにリポジトリを clone します。ただし、ラボのコードはこれらを使いません。書籍のスクリプトをそのまま流用しているために含まれています。

クラスタは約 4 分で起動しました。イメージバージョンは 2.2.87-debian12、ゾーンは自動で us-central1-a が選ばれました。

図2 クラスタの構成。赤枠:マスターノードとワーカーノード 2 台のマシンタイプ e2-standard-4
SS-02b_nodes.png

図3 クラスタの構成(続き)。赤枠:オプション コンポーネント JUPYTER
SS-02c_components.png

なぜ --public-ip-address が必要なのか

イメージバージョン 2.2 以降のクラスタは、デフォルトで内部 IP アドレスだけを持つ VM として作成されます。Google の API には限定公開の Google アクセスで到達できますが、インターネットには出られません。

ラボの初期化アクションは、GitHub からの git clone や apt-get を実行します。そのため、外部 IP(または Cloud NAT)がないと失敗します。ラボが --public-ip-address を追加しているのはこのためです。

組織のポリシー(constraints/compute.vmExternalIpAccess)で外部 IP が禁止されているプロジェクトでは、--public-ip-address を外し、初期化アクションも外してください。前述のとおり、ラボのコードは初期化アクションを使いません。

詳細:公式ドキュメント(Managed Service for Apache Spark クラスタのネットワーク構成)、公式ドキュメント(初期化アクション)


コードの実行方法

ラボでは、クラスタの ウェブ インターフェース タブから JupyterLab を開き、ノートブックでコードを実行します。JupyterLab はコンポーネント ゲートウェイ経由で公開されます。

図4 ウェブ インターフェース タブ。コンポーネント ゲートウェイ経由で開ける画面の一覧。赤枠:JupyterLab
SS-03.png

今回の検証環境では、ブラウザから JupyterLab を開けなかったため、ノートブックのセルを 1 つの Python ファイルにまとめ、PySpark ジョブとして送信しました。

gcloud dataproc jobs submit pyspark lab_gsp271.py \
  --cluster=ch6cluster \
  --region="${REGION}"

コードはラボのセルと同じです。変更したのは次の 2 点だけです。

  • ノートブック専用の構文 PROJECT=!gcloud config get-value project を、subprocess での取得に置き換えた
  • グラフを画面に表示する代わりに、PNG ファイルとして Cloud Storage に保存した

ジョブの出力は、ターミナルと、コンソールの ジョブの詳細 画面の両方で確認できます。以降の図は、ジョブの詳細画面の出力です。ジョブは約 80 秒で完了しました。

詳細:公式ドキュメント(コンポーネント ゲートウェイ)、gcloud dataproc jobs submit pyspark(英語)


Spark セッションを作成する

import subprocess
import os

PROJECT = subprocess.check_output(['gcloud', 'config', 'get-value', 'project'], text=True).strip()
BUCKET = PROJECT + '-dsongcp'

from pyspark.sql import SparkSession
from pyspark import SparkContext
sc = SparkContext('local', 'logistic')
spark = SparkSession \
    .builder \
    .appName("Logistic regression w/ Spark ML") \
    .getOrCreate()

from pyspark.mllib.classification import LogisticRegressionWithLBFGS
from pyspark.mllib.regression import LabeledPoint

ここで SparkContext('local', …) に注目してください。local を指定すると、Spark はマスターノードのプロセスの中だけで動き、ワーカーノードは使われません。実行時に確認すると、次のとおりでした。

sc.master = local / defaultParallelism = 1 / Spark 3.5.3

1 万数千行のデータであれば local でも十分な速さで処理できます。ただし、本来クラスタで分散処理したい場合は、SparkContext の master を指定せず、YARN に任せる必要があります。

また、pyspark.mllib は RDD ベースの古い API で、現在はメンテナンスモードです。新しく開発する場合は、DataFrame ベースの pyspark.ml を使ってください。

詳細:Spark 3.5 MLlib の線形手法:ロジスティック回帰(英語)


データを読み込んでクリーニングする

学習日と評価日

trainday.csv は、日付ごとに学習用(True)か評価用(False)かを示す表です。

traindays = spark.read \
    .option("header", "true") \
    .csv('gs://{}/data-science-on-gcp/flights/trainday.csv'.format(BUCKET))
traindays.createOrReplaceTempView('traindays')

spark.sql("SELECT * from traindays LIMIT 5").show()

図5 Spark セッションと trainday.csv の読み込み結果。赤枠:local で動いていることを示す出力と、学習日の区分
SS-05_traindays.png

フライトデータの列を対応付ける

フライトデータの CSV にはヘッダ行がありません。Spark は列を _c0、_c1… と名付けるため、列番号で必要な列を取り出します。

from pyspark.sql.functions import col, when

inputs = 'gs://{}/data-science-on-gcp/flights/tzcorr/all_flights-00001-*'.format(BUCKET)
raw_flights = spark.read.option("header", "false").csv(inputs)

flights = raw_flights.select(
    col("_c0").alias("FL_DATE"),
    col("_c15").cast("double").alias("DEP_DELAY"),
    col("_c16").cast("double").alias("TAXI_OUT"),
    col("_c22").cast("double").alias("ARR_DELAY"),
    when(col("_c23") == "1.00", "True").otherwise("False").alias("CANCELLED"),
    when(col("_c24") == "1.00", "True").otherwise("False").alias("DIVERTED"),
    col("_c26").cast("double").alias("DISTANCE")
)
flights.createOrReplaceTempView('flights')

trainquery = """
SELECT
DEP_DELAY, TAXI_OUT, ARR_DELAY, DISTANCE
FROM flights f
JOIN traindays t ON f.FL_DATE = t.FL_DATE
WHERE
t.is_train_day = 'True' AND
f.CANCELLED = 'False' AND
f.DIVERTED = 'False' AND
f.DEP_DELAY IS NOT NULL AND
f.ARR_DELAY IS NOT NULL
"""
traindata = spark.sql(trainquery)
traindata.describe().show()

欠航(CANCELLED)や目的地変更(DIVERTED)のフライトには到着遅延がないため、除外しています。

図6 学習データの統計量。赤枠:4 列とも件数は 12,303
SS-06_describe.png

ラボの列対応には1か所誤りがある

tzcorr.json の列定義と照らし合わせると、次のようになっています。

列 実際の列名 ラボでの扱い
_c23 CANCELLED CANCELLED
_c24 CANCELLATION_CODE(A/B/C または空) DIVERTED
_c25 DIVERTED (使われていない)
_c26 DISTANCE DISTANCE

ラボは _c24 を DIVERTED として扱っていますが、実際は欠航理由のコードです。"1.00" になることはないため、DIVERTED は常に False になります。

**ただし、今回のデータでは結果に影響しません。**目的地変更になったフライトはすべて到着遅延が空で、ARR_DELAY IS NOT NULL の条件で除外されるためです(学習用シャードでは 27 件すべて)。コードを流用する場合は、_c25 に直してください。


ロジスティック回帰モデルを学習する

各行を LabeledPoint(正解ラベルと特徴量の組)に変換し、学習します。

def to_example(fields):
    return LabeledPoint(\
              float(fields['ARR_DELAY'] < 15), #ontime? \
              [ \
                  fields['DEP_DELAY'], \
                  fields['TAXI_OUT'],  \
                  fields['DISTANCE'],  \
              ])

examples = traindata.rdd.map(to_example)
lrmodel = LogisticRegressionWithLBFGS.train(examples, intercept=True)
print(lrmodel.weights,lrmodel.intercept)

intercept=True は、入力がすべて 0 のときの予測が 0 にならないことを表します。

学習データ 12,303 件のうち、定刻(到着遅延 15 分未満)は 10,650 件(86.6%)でした。定刻のほうが大幅に多い、偏りのあるデータです。

重みと予測

学習後の重みと切片は次のとおりです。

[-0.16903723994894598,-0.13173260912370688,0.00015285415780823042] 5.556528108526656

出発遅延とタキシングアウト時間の重みは負で、長いほど定刻の確率が下がります。飛行距離の重みはわずかに正です。長距離便は飛行中に遅れを取り戻しやすいと解釈できます。

ラボの記載([-0.179, -0.135, 0.00048] 5.40)とは値が異なります。ラボの説明文は、より大きなデータで作られたと考えられます。符号と桁は一致しています。

print(lrmodel.predict([6.0,12.0,594.0]))    # 出発遅延 6 分
print(lrmodel.predict([36.0,12.0,594.0]))   # 出発遅延 36 分

結果は 1(定刻)と 0(遅延)です。

確率としきい値

predict はデフォルトでしきい値 0.5 を使い、0 か 1 を返します。clearThreshold() でしきい値を外すと、確率が返ります。

lrmodel.clearThreshold()
print(lrmodel.predict([6.0,12.0,594.0]))
print(lrmodel.predict([36.0,12.0,594.0]))

lrmodel.setThreshold(0.7)
print(lrmodel.predict([6.0,12.0,594.0]))
print(lrmodel.predict([36.0,12.0,594.0]))

ラボの設定は「定刻に着く確率が 70% を下回ったら、到着後の会議をキャンセルする」というものです。そのため、しきい値を 0.7 にしています。

図7 学習結果。赤枠:重みと切片、しきい値 0.5 での予測(上)。確率と、しきい値 0.7 での予測(下)
SS-07-08_weights_prob.png

出発遅延 6 分の便は定刻の確率が 0.955、36 分の便は 0.117 でした。


モデルを保存して復元する

MLlib のモデルは、Cloud Storage に直接保存できます。

MODEL_FILE='gs://' + BUCKET + '/flights/sparkmloutput/model'
os.system('gcloud storage rm -r ' + MODEL_FILE)
lrmodel.save(sc, MODEL_FILE)

lrmodel = 0

from pyspark.mllib.classification import LogisticRegressionModel
lrmodel = LogisticRegressionModel.load(sc, MODEL_FILE)
lrmodel.setThreshold(0.7)

保存先に既存のファイルがあると保存に失敗するため、先に削除しています。初回は削除対象がなく、gcloud storage rm はエラーを返しますが、問題ありません。

図8 保存と復元。赤枠:復元後の重みは保存前と同じ。出発遅延 36 分の便は 0(遅延)、8 分の便は 1(定刻)
SS-07b_restore_predict.png

図9 保存されたモデル。赤枠:メタデータ(metadata/)と重み(data/、Parquet 形式)
SS-09_model.png


モデルの挙動を確認する

しきい値を外し、入力を 1 つずつ変えたときの確率をグラフにします。

lrmodel.clearThreshold()
dist = np.arange(10, 2000, 10)
prob = [lrmodel.predict([20, 10, d]) for d in dist]
plt.plot(dist, prob)
plt.xlabel('distance (miles)')
plt.ylabel('probability of ontime arrival')

図10 飛行距離と定刻の確率(出発遅延 20 分、タキシングアウト 10 分)
SS-10_distance.png

距離が 10 マイルから 1,990 マイルに伸びても、確率は 0.70 から 0.76 に上がるだけです。距離の影響は小さいことがわかります。

delay = np.arange(-20, 60, 1)
prob = [lrmodel.predict([d, 10, 500]) for d in delay]
plt.plot(delay, prob)
plt.xlabel('departure delay (minutes)')
plt.ylabel('probability of ontime arrival')

図11 出発遅延と定刻の確率(タキシングアウト 10 分、距離 500 マイル)
SS-11_delay.png

出発遅延の影響は大きく、確率は 10 分を過ぎたあたりから急に下がり、約 20 分で 0.7 を下回ります(10 分で 0.93、20 分で 0.72、30 分で 0.32)。-20 分では 0.9995、59 分では 0.0035 でした。


モデルを評価する

評価には、学習に使っていないシャード(all_flights-00002-*)の評価日のデータを使います。

inputs = 'gs://{}/data-science-on-gcp/flights/tzcorr/all_flights-00002-*'.format(BUCKET)
raw_test_flights = spark.read.option("header", "false").csv(inputs)
# 列の対応付けは学習時と同じ(省略)

testquery = trainquery.replace("t.is_train_day = 'True'", "t.is_train_day = 'False'")
testdata = spark.sql(testquery)
examples = testdata.rdd.map(to_example)
print(f"Test examples count: {examples.count()}")

評価データは 12,307 件でした。ラボの記載(82,184 件)とは異なります。重みと同様に、ラボの説明文は別のデータで作られたものと考えられます。

予測確率が 0.7 未満のフライトを「キャンセルする」、0.7 以上を「キャンセルしない」に分け、それぞれの正解率を求めます。

def eval(labelpred):
    cancel = labelpred.filter(lambda data: data[1] < 0.7)
    nocancel = labelpred.filter(lambda data: data[1] >= 0.7)
    corr_cancel = cancel.filter(lambda data: data[0] == int(data[1] >= 0.7)).count()
    corr_nocancel = nocancel.filter(lambda data: data[0] == int(data[1] >= 0.7)).count()
    cancel_denom = cancel.count()
    nocancel_denom = nocancel.count()
    if cancel_denom == 0:
        cancel_denom = 1
    if nocancel_denom == 0:
        nocancel_denom = 1
    return {'total_cancel': cancel.count(),
            'correct_cancel': float(corr_cancel)/cancel_denom,
            'total_noncancel': nocancel.count(),
            'correct_noncancel': float(corr_nocancel)/nocancel_denom}

lrmodel.clearThreshold()
labelpred = examples.map(lambda p: (p.label, lrmodel.predict(p.features)))
print(eval(labelpred))

labelpred = labelpred.filter(lambda data: data[1] > 0.65 and data[1] < 0.75)
print(eval(labelpred))

図12 評価結果。赤枠:全フライト(上)と、確率が 0.65〜0.75 のフライト(下)
SS-12_eval.png

対象 キャンセルする(件数/正解率) キャンセルしない(件数/正解率)
全フライト 1,515 件/78.7% 10,792 件/96.1%
確率 0.65〜0.75 91 件/39.6% 101 件/64.4%

全体では高い正解率ですが、**しきい値の近くでは大きく下がります。**確率が 0.7 前後のフライトは、モデルにとっても判断が難しいということです。

ここでの「キャンセル」は、欠航ではなく「会議をキャンセルする判断」を指します。欠航したフライトは、学習データと評価データの両方から除外済みです。


落とし穴

Spark が local モードで動いている

ラボのコードは SparkContext('local', …) で Spark を起動するため、マスターノードの中だけで処理が完結し、ワーカーノードは使われません。クラスタとして分散処理させたい場合は、master を指定しないでください。

ラボの列対応に誤りがある

_c24 は DIVERTED ではなく CANCELLATION_CODE です。今回のデータでは結果に影響しませんが、コードを流用する場合は _c25 に直してください。

ラボの説明文の数値は現在のデータと合わない

重み、評価データの件数、正解率は、いずれもラボの説明文と一致しませんでした。現在のシャードで計算すると、学習データ 12,303 件、評価データ 12,307 件です。符号や傾向は一致するため、手順の確認には使えます。

イメージ 2.2 以降は内部 IP のみがデフォルト

外部 IP がないと、初期化アクションの git clone や apt-get が失敗します。外部 IP が使えない環境では、初期化アクションを外してください。ラボのコードは初期化アクションを使いません。

中断したつもりの作成コマンドが実行されている

gcloud dataproc clusters create を途中で中断しても、リクエストが API に届いていれば、クラスタの作成は続きます。今回、再実行したところ ALREADY_EXISTS になりました。再実行の前に、gcloud dataproc clusters list で確認してください。

ステージングバケットと一時バケットが残る

--bucket と --temp-bucket を指定しないと、クラスタの作成時に dataproc-staging-<リージョン>-<プロジェクト番号>-<ランダム文字列> と dataproc-temp-… の 2 つのバケットが自動で作られます。**クラスタを削除してもバケットは残ります。**一時バケットには 90 日の TTL がありますが、ステージングバケットにはありません。

詳細:公式ドキュメント(ステージング バケットと一時バケット)


費用について

費用の大半はクラスタです。クラスタは、ジョブを実行していなくても、存在している間ずっと課金されます。

項目 内容
VM e2-standard-4 × 3 台(12 vCPU)
Managed Service for Apache Spark の料金 vCPU 数と稼働時間に応じて課金
ディスク 500 GB の標準永続ディスク × 3 台

今回はクラスタの作成から削除まで約 16 分でした。作業が終わったら、すぐにクラスタを削除してください。

詳細:公式ドキュメント(Managed Service for Apache Spark の料金)


クリーンアップ

gcloud dataproc clusters delete ch6cluster --region="${REGION}" --quiet

gcloud storage rm -r "gs://${BUCKET}"

gcloud storage ls | grep -E "dataproc-(staging|temp)-${REGION}"

最後のコマンドで、自動で作られたステージングバケットと一時バケットの名前を確認します。同じリージョンに他のクラスタがないことを確認してから、表示された 2 つのバケットを gcloud storage rm -r で削除してください。複数のクラスタで共有されている場合があります。

gcloud projects remove-iam-policy-binding "${PROJECT_ID}" \
  --member="serviceAccount:${SA_EMAIL}" \
  --role="roles/dataproc.worker" \
  --condition=None

gcloud iam service-accounts delete "${SA_EMAIL}" --quiet

IAM のバインディングは、サービスアカウントを削除する前に外してください。先にサービスアカウントを削除すると、バインディングが deleted:serviceAccount:… として残ります。


クラスタの管理を任せて、Spark のコードに集中する。 Managed Service for Apache Spark では、数分で起動したクラスタの上で、MLlib のモデルの学習、保存、評価までを実行できます。一方で、local モードのように、クラスタがあってもコードの書き方次第で分散処理されない場合があります。サンプルコードを使うときは、実行モード、列の対応、クリーンアップまでを確認してから使ってください。

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?