I'm new to distributed ML and currently doing my personal project
I train my model using PySpark on Cloud Dataproc and build the pipeline like the code below
spark = SparkSession.builder.appName('sparkify-train').getOrCreate()
df = spark.read.parquet(path)
gbt = GBTClassifier()
paramGrid = ParamGridBuilder() \
.addGrid(gbt.maxDepth, [4,8,12]) \
.addGrid(gbt.maxIter, [5,10,15]) \
.addGrid(gbt.featuresCol, ["features_1", "features_2"]) \
.build()
tvs = TrainValidationSplit(estimator=gbt,
estimatorParamMaps=paramGrid,
evaluator=BinaryClassificationEvaluator(),
trainRatio=0.75)
tvs.fit(df)
Note: I don't show feature extraction and other stuff on the code above
But on CPU utilization chart, the master node is the only busy node instead of three workers node
I create Dataproc cluster using cloud shell with command below. This command is generated from Dataproc create cluster UI.
gcloud beta dataproc clusters create sparkify --enable-component-gateway --region asia-southeast2 --subnet default --zone asia-southeast2-c
--master-machine-type custom-1-3840 --master-boot-disk-type pd-ssd --master-boot-disk-size 50 --num-workers 3 --worker-machine-type custom-1-3840 --worker-boot-disk-type pd-ss
d --worker-boot-disk-size 30 --image-version 2.0-ubuntu18 --optional-components JUPYTER,ZEPPELIN --project wskt-trek
and this is the cluster properties from command above:

I just wondering is there anything wrong in my code? Or are there any settings I miss? I thought all of the workers CPU should have high utilization. Thanks!
