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/

クラスタを作成する
ラボの変更点を反映したコマンド
ラボでは、リポジトリの 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

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

なぜ --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

今回の検証環境では、ブラウザから 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 で動いていることを示す出力と、学習日の区分

フライトデータの列を対応付ける
フライトデータの 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

ラボの列対応には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 での予測(下)

出発遅延 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(定刻)

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

モデルの挙動を確認する
しきい値を外し、入力を 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 分)

距離が 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 マイル)

出発遅延の影響は大きく、確率は 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 のフライト(下)

| 対象 | キャンセルする(件数/正解率) | キャンセルしない(件数/正解率) |
|---|---|---|
| 全フライト | 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 モードのように、クラスタがあってもコードの書き方次第で分散処理されない場合があります。サンプルコードを使うときは、実行モード、列の対応、クリーンアップまでを確認してから使ってください。