Sie können einen Job über eine jobs.submit API als HTTP- oder programmatische Anfrage an einen vorhandenen Managed Service for Apache Spark-Cluster übergeben. Verwenden Sie dazu das gcloud-Befehlszeilentool der Google Cloud CLI in einem lokalen Terminalfenster oder in Cloud Shell oder aus der Google Cloud -Konsole in einem lokalen Browser. Sie können auch eine SSH-Verbindung zur Masterinstanz in Ihrem Cluster herstellen und dann einen Job direkt aus der Instanz ausführen, ohne Managed Service for Apache Spark zu verwenden.
Job-Nebenläufigkeit : Sie können die maximale Anzahl gleichzeitiger Managed Service for Apache Spark-Jobs mit dem Clusterattribut dataproc:dataproc.scheduler.max-concurrent-jobs konfigurieren, wenn Sie einen Cluster erstellen. Wenn dieser Attributwert nicht festgelegt ist, wird die Obergrenze für gleichzeitige Jobs als max((masterMemoryMb - 3584) / masterMemoryMbPerJob, 5) berechnet.
masterMemoryMb wird vom Maschinentyp der Master-VM bestimmt.
masterMemoryMbPerJob ist standardmäßig 1024, kann aber bei der Clustererstellung mit dem Clusterattribut dataproc:dataproc.scheduler.driver-size-mb konfiguriert werden.
Job einreichen
Classpath-Isolation und benutzerdefinierte JARs:Kopieren Sie keine benutzerdefinierten „Fat JARs“ oder Bundle-JARs (z. B. Apache Iceberg-Laufzeit oder Google Cloud -Bundles) direkt in Cluster-Systemverzeichnisse wie /usr/lib/spark/jars/. Wenn Sie benutzerdefinierte JARs in Systemverzeichnisse einfügen, wird der Klassenpfad des Managed Service for Apache Spark-Agents mit transitiven Abhängigkeiten (z. B. Guava- oder Hadoop-Bibliotheken) überladen, die mit den integrierten Bibliotheken des Agents in Konflikt stehen. Dieser Konflikt kann zu Fehlern bei der Klassenauflösung führen (z. B. ClassNotFoundException oder NoClassDefFoundError), wenn Jobverwaltungsoperationen wie das Abbrechen eines Jobs ausgeführt werden. Dies kann dazu führen, dass der Agent abstürzt und YARN-Anwendungen verwaist sind.
Verwenden Sie stattdessen einen der folgenden unterstützten Ansätze:
- Jobspezifische Abhängigkeiten:Geben Sie den Cloud Storage-Pfad zu Ihren JAR-Dateien an, wenn Sie den Job mit dem Flag
--jarsin der Google Cloud CLI, dem Feld JAR-Dateien in der Google Cloud Console oder dem FeldjarFileUrisin der API einreichen. Spark verteilt die Abhängigkeiten an den Treiber und die Executors für den Job, ohne den Klassenpfad des Agents zu beeinträchtigen. - Clusterweite Abhängigkeiten:Geben Sie Abhängigkeiten bei der Clustererstellung mit dem
spark:spark.jars-Clusterattribut an: Dadurch wird Managed Service for Apache Spark angewiesen,--properties="spark:spark.jars=gs://YOUR_BUCKET/jar-1.jar,gs://YOUR_BUCKET/jar-2.jar"
/etc/spark/conf/spark-defaults.confzu konfigurieren. Die Abhängigkeiten werden automatisch auf die Spark-Treiber- und Executor-Klassenpfade für alle Jobs verteilt, während der knotenlokale Managed Service for Apache Spark-Agent sauber isoliert bleibt.
Console
Öffnen Sie die Seite Job senden für Managed Service for Apache Spark in der Google Cloud Console in Ihrem Browser.
Spark-Job – Beispiel
Zum Senden eines Spark-Beispieljobs füllen Sie die Felder auf der Seite Job senden so aus:
- Wählen Sie den Namen des Clusters aus der Clusterliste aus.
- Legen Sie für Job type (Jobtyp) den Wert
Sparkfest. - Legen Sie für Main class or jar (Hauptklasse oder JAR-Datei) den Wert
org.apache.spark.examples.SparkPifest. - Legen Sie für Arguments (Argumente) das einzelne Argument
1000fest. - Fügen Sie
file:///usr/lib/spark/examples/jars/spark-examples.jarzu Jar-Dateien (oder dem API-FeldjarFileUris) hinzu:file:///gibt ein Hadoop LocalFileSystem-Schema an./usr/lib/spark/examples/jars/spark-examples.jarwurde von Managed Service for Apache Spark auf dem Masterknoten des Clusters installiert, als der Cluster erstellt wurde. Dieser Pfad wird nur für vorinstallierte Beispiele verwendet, die von Managed Service for Apache Spark bereitgestellt werden. Kopieren Sie keine benutzerdefinierten JARs in Systemverzeichnisse (siehe den Hinweis unter Job senden).- Alternativ können Sie einen Cloud Storage-Pfad (
gs://your-bucket/your-jarfile.jar) oder einen Hadoop Distributed File System-Pfad (hdfs://path-to-jar.jar) zu einer Ihrer JAR-Dateien angeben. Wenn Sie den Job über die API einreichen, geben Sie diesen Pfad im FeldjarFileUrisan.
Klicken Sie auf Submit (Senden), um den Job zu starten. Nach dem Start wird der Job der Jobliste hinzugefügt.
Klicken Sie auf die Job-ID, um die Seite Jobs zu öffnen, auf der Sie die Treiberausgabe des Jobs anzeigen können. Da die generierten Ausgabezeilen die Breite des Browserfensters überschreiten, klicken Sie auf das Kästchen Zeilenumbruch, um den gesamten Ausgabetext in der Ansicht darzustellen und das berechnete Ergebnis für pi einzublenden.
Sie können die Treiberausgabe des Jobs über die Befehlszeile mit dem unten gezeigten Befehl gcloud dataproc jobs wait aufrufen. Weitere Informationen finden Sie unter Jobausgabe ansehen.
Kopieren Sie die Projekt-ID und fügen Sie sie als den Wert für das Flag --project ein. Kopieren Sie anschließend die Job-ID (in der Jobliste angezeigt) und fügen Sie sie als endgültiges Argument ein.
gcloud dataproc jobs wait job-id \ --project=project-id \ --region=region
Hier sind Snippets aus der Treiberausgabe für den Beispieljob SparkPi:
... 2015-06-25 23:27:23,810 INFO [dag-scheduler-event-loop] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Stage 0 (reduce at SparkPi.scala:35) finished in 21.169 s 2015-06-25 23:27:23,810 INFO [task-result-getter-3] cluster.YarnScheduler (Logging.scala:logInfo(59)) - Removed TaskSet 0.0, whose tasks have all completed, from pool 2015-06-25 23:27:23,819 INFO [main] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Job 0 finished: reduce at SparkPi.scala:35, took 21.674931 s Pi is roughly 3.14189648 ... Job [c556b47a-4b46-4a94-9ba2-2dcee31167b2] finished successfully. driverOutputUri: gs://sample-staging-bucket/google-cloud-dataproc-metainfo/cfeaa033-749e-48b9-... ...
gcloud
Zum Senden eines Jobs an einen Managed Service for Apache Spark-Cluster führen Sie den gcloud CLI-Befehl gcloud dataproc jobs submit lokal in einem Terminalfenster oder in Cloud Shell aus.
gcloud dataproc jobs submit job-command \ --cluster=cluster-name \ --region=region \ other dataproc-flags \ -- job-args
- Listen Sie das öffentlich zugängliche
hello-world.pyin Cloud Storage auf. Dateiliste:gcloud storage cat gs://dataproc-examples/pyspark/hello-world/hello-world.py
#!/usr/bin/python import pyspark sc = pyspark.SparkContext() rdd = sc.parallelize(['Hello,', 'world!']) words = sorted(rdd.collect()) print(words)
- Senden Sie den PySpark-Job an Managed Service for Apache Spark.
Terminalausgabe:gcloud dataproc jobs submit pyspark \ gs://dataproc-examples/pyspark/hello-world/hello-world.py \ --cluster=cluster-name \ --region=region
Waiting for job output... … ['Hello,', 'world!'] Job finished successfully.
- Führen Sie das vorinstallierte SparkPi-Beispiel auf dem Masterknoten des Managed Service for Apache Spark-Clusters aus. Der Pfad
file:///usr/lib/spark/examples/jars/spark-examples.jarist nur für vorinstallierte Beispiele vorgesehen. Informationen zu benutzerdefinierten JAR-Abhängigkeiten finden Sie im Warnhinweis unter Job senden. Terminalausgabe:gcloud dataproc jobs submit spark \ --cluster=cluster-name \ --region=region \ --class=org.apache.spark.examples.SparkPi \ --jars=file:///usr/lib/spark/examples/jars/spark-examples.jar \ -- 1000
Job [54825071-ae28-4c5b-85a5-58fae6a597d6] submitted. Waiting for job output… … Pi is roughly 3.14177148 … Job finished successfully. …
REST
In diesem Abschnitt wird gezeigt, wie Sie einen Spark-Job senden, um den ungefähren Wert von pi mithilfe der Managed Service for Apache Spark-API jobs.submit zu berechnen.
Ersetzen Sie diese Werte in den folgenden Anfragedaten:
- project-id: Google Cloud Projekt-ID
- region: Cluster-Region
- clusterName: Clustername
HTTP-Methode und URL:
POST https://dataproc.googleapis.com/v1/projects/project-id/regions/region/jobs:submit
JSON-Text anfordern:
{
"job": {
"placement": {
"clusterName": "cluster-name"
},
"sparkJob": {
"args": [
"1000"
],
"mainClass": "org.apache.spark.examples.SparkPi",
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
}
}
}
Wenn Sie die Anfrage senden möchten, maximieren Sie eine der folgenden Optionen:
Sie sollten eine JSON-Antwort ähnlich wie diese erhalten:
{
"reference": {
"projectId": "project-id",
"jobId": "job-id"
},
"placement": {
"clusterName": "cluster-name",
"clusterUuid": "cluster-Uuid"
},
"sparkJob": {
"mainClass": "org.apache.spark.examples.SparkPi",
"args": [
"1000"
],
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
},
"status": {
"state": "PENDING",
"stateStartTime": "2020-10-07T20:16:21.759Z"
},
"jobUuid": "job-Uuid"
}
Java
- Clientbibliothek installieren
- Standardanmeldedaten für Anwendungen einrichten
- Führen Sie den Code aus
Python
- Clientbibliothek installieren
- Standardanmeldedaten für Anwendungen einrichten
- Führen Sie den Code aus
Go
Node.js
Job direkt an Cluster senden
Wenn Sie einen Job direkt auf Ihrem Cluster ohne den Managed Service for Apache Spark ausführen möchten, stellen Sie eine SSH-Verbindung zum Masterknoten Ihres Clusters her und führen Sie den Job dann auf dem Masterknoten aus.
Nachdem Sie eine SSH-Verbindung zur VM-Masterinstanz hergestellt haben, führen Sie die folgenden Schritte in einem Terminalfenster im Masterknoten des Clusters aus:
- Öffnen Sie eine Spark-Shell.
- Führen Sie einen Spark-Job aus, um die Anzahl der Zeilen in einer (siebenzeiligen) Python-Datei „hello-world“ zu zählen, die sich in einer öffentlich zugänglichen Cloud Storage-Datei befindet.
Beenden Sie die Shell.
user@cluster-name-m:~$ spark-shell ... scala> sc.textFile("gs://dataproc-examples" + "/pyspark/hello-world/hello-world.py").count ... res0: Long = 7 scala> :quit
Bash-Jobs in Managed Service for Apache Spark ausführen
Möglicherweise möchten Sie ein Bash-Skript als Managed Service for Apache Spark-Job ausführen, weil die von Ihnen verwendeten Engines nicht als Managed Service for Apache Spark-Job auf oberster Ebene unterstützt werden oder weil Sie vor dem Start eines Jobs mit hadoop oder spark-submit aus dem Skript zusätzliche Einrichtungsschritte oder Berechnung von Argumenten vornehmen müssen.
Pig-Beispiel
Angenommen, Sie haben ein hello.sh-Bash-Skript in Cloud Storage kopiert:
gcloud storage cp hello.sh gs://${BUCKET}/hello.shDa der Befehl pig fs Hadoop-Pfade verwendet, kopieren Sie das Skript aus Cloud Storage an ein Ziel als file:///, damit es sich im lokalen Dateisystem statt in HDFS befindet. Die nachfolgenden sh-Befehle verweisen automatisch auf das lokale Dateisystem und erfordern nicht das Präfix file:///.
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
-e='fs -cp -f gs://${BUCKET}/hello.sh file:///tmp/hello.sh; sh chmod 750 /tmp/hello.sh; sh /tmp/hello.sh'Alternativ können Sie das Cloud Storage-Shell-Skript als Argument --jars angeben, da das Argument --jars für die Übermittlung von Managed Service for Apache Spark-Jobs eine Datei in ein temporäres Verzeichnis stellt, das für die Lebensdauer des Jobs erstellt wurde:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=gs://${BUCKET}/hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'Beachten Sie, dass das Argument --jars auch auf ein lokales Skript verweisen kann:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'