Spark-Anwendungen interaktiv ausführen
Sie können Ihre Spark-Anwendungen unter Verwendung der Kernel-API interaktiv ausführen.
Kernel as a Service bietet:
- Jupyter-Kernel als Entität der ersten Klasse
- Ein dedizierter Cluster pro Kernel
- Cluster und Kernel angepasst mit Benutzerbibliotheken
Kernel-API verwenden
Sie können eine Spark-Anwendung interaktiv mithilfe der Kernel-API ausführen. Jede Anwendung wird in einem Kernel in einem dedizierten Cluster ausgeführt. Alle Konfigurationseinstellungen, die Sie über die API übergeben, überschreiben die Standardkonfigurationen.
Gehen Sie wie folgt vor, um eine Spark-Anwendung interaktiv auszuführen:
Generieren Sie ein Zugriffstoken, wenn Sie kein Token haben. Siehe „Zugriffstoken generieren “.
Das Token in eine Variable exportieren:
export TOKEN=<token generated>Eine Instanz von „ Analytics Engine powered by Apache Spark “ bereitstellen. Sie benötigen die Rolle Administrator im Projekt oder im Bereitstellungsbereich, um eine Instanz bereitzustellen. Siehe „Bereitstellung einer Instanz “.
Rufen Sie den Spark-Kernelendpunkt für die Instanz ab.
- Klicken Sie im Navigationsmenü
in Cloud Pak for Dataauf Services > Instanzen, suchen Sie die Instanz und klicken Sie darauf, um die Instanzdetails anzuzeigen.
- Kopieren und speichern Sie unter "Zugriffsinformationen" den Spark-Kernelendpunkt.
- Klicken Sie im Navigationsmenü
Erstellen Sie einen Kernel mit dem generierten Kernelendpunkt und Zugriffstoken. Dieses Beispiel enthält die obligatorischen Mindestparameter, die erforderlich sind:
curl -k -X POST <KERNEL_ENDPOINT> -H "Authorization: Bearer ${TOKEN}" -d '{"name":"scala" }'
Kernel-JSON und unterstützte Parameter erstellen
Wenn Sie den Kernel-Service verwenden, um die REST-API zum Erstellen von Kernels zu starten, können Sie den Nutzdaten erweiterte Konfigurationen hinzufügen.
Sehen Sie sich die folgenden Beispielnutzdaten und Parameterlisten an, die Sie in der API "create kernel" übergeben können:
curl -k -X POST <KERNEL_ENDPOINT> -H "Authorization: Bearer ${TOKEN}" -d '{
"name": "scala",
"kernel_size": {
"cpu": 1,
"memory": "1g"
},
"engine": {
"type": "spark",
"conf": {
"spark.ui.reverseProxy": "false",
"spark.eventLog.enabled": "false"
},
"size": {
"num_workers": "2",
"worker_size": {
"cpu": 1,
"memory": "1g"
}
}
}
}'
Antwort:
{
"id": "<kernel_id>",
"name": "scala",
"connections": 0,
"last_activity": "2021-07-16T11:06:26.266275Z",
"execution_state": "starting"
}
Spark-Kernel-API unter Verwendung einer angepassten Spark-Laufzeit
Beispiel für Eingabenutzdaten zum Ändern der Spark-Laufzeitversion:
{
"name": "python310",
"kernel_size": {
"cpu": 1,
"memory": "1g"
},
"engine": {
"type": "spark",
"runtime": {
"spark_version": "3.4"
},
"size": {
"num_workers": "2",
"worker_size": {
"cpu": 1,
"memory": "1g"
}
}
}
}
Parameter der Spark-Kernel-API
Sie können die folgenden Parameter in der Kernel-API verwenden:
| Ihren Namen | Erforderlich/Optional | Typ | Beschreibung |
|---|---|---|---|
| Name | Erforderlich | Zeichenfolge | Gibt den Kernelnamen an. Unterstützte Werte sind scala, r, r42, python39, python310 |
| Maschine | Optionale | Schlüssel-Wert-Paare | Gibt die Spark-Laufzeit mit Konfigurations-und Versionsinformationen an |
| engine.runtime.spark_version | Optionale | Zeichenfolge | Gibt die für den Kernel zu verwendende Spark-Laufzeitversion an. IBM Cloud Pak for Data unterstützt Spark- 3.4 |
| Typ | Erforderlich, wenn die Engine angegeben ist | Zeichenfolge | Gibt den Laufzeittyp des Kernels an Derzeit wird nur spark unterstützt. |
| conf | Optionale | JSON-Objekt mit Schlüssel-Wert-Paaren | Gibt die Spark-Konfigurationswerte an, die die vordefinierten Werte überschreiben. |
| Umgebung | Optionale | JSON-Objekt mit Schlüssel-Wert-Paaren | Gibt die für den Job erforderlichen Spark-Umgebungsvariablen an |
| Größe | Optionale | Akzeptiert die Parameter num_workers, worker_size und master_size |
|
| Anzahl der Mitarbeiter | Erforderlich, wenn Größe angegeben wird | Ganzzahl | Gibt die Anzahl der Worker-Knoten im Spark-Cluster an. num_workers entspricht der Anzahl der gewünschten Ausführenden. Standardmäßig wird pro Worker-Knoten ein Executor verwendet. Es werden maximal 50 Ausführende unterstützt. |
| Arbeitnehmergröße | Optionale | memoryNimmt die Parameter cpu und an. |
|
| cpu | Erforderlich, wenn worker_size angegeben ist |
Ganzzahl | Legt die CPU-Leistung für den Worker-Knoten fest. Standardmäßig ist 1 CPU eingestellt. Maximal 10 CPUs |
| Speicher | Erforderlich, wenn worker_size angegeben ist |
Ganzzahl | Legt die Speichermenge für jeden Worker-Knoten fest. Die Standardeinstellung beträgt 1 GB. Maximal 40 GB |
| Mastergröße | Optionale | Akzeptiert die Parameter cpu und memory |
|
| cpu | Erforderlich, wenn master_size angegeben ist |
Ganzzahl | Gibt die CPU-Kapazität für den Masterknoten an. Standardmäßig ist 1 CPU eingestellt. Maximal 10 CPUs |
| Speicher | Erforderlich, wenn master_size angegeben ist |
Ganzzahl | Gibt die Speicherkapazität für jeden Masterknoten an. Die Standardeinstellung beträgt 1 GB. Maximal 40 GB |
| kernelgröße | Optionale | Akzeptiert die Parameter cpu und memory |
|
| cpu | Erforderlich, wenn kernel_size angegeben ist |
Ganzzahl | Gibt die CPU-Kapazität für den Kernelknoten an Standardmäßig ist 1 CPU eingestellt. Maximal 10 CPUs |
| Speicher | Erforderlich, wenn kernel_size angegeben ist |
Ganzzahl | Gibt die Speicherkapazität für jeden Kernelknoten an. Die Standardeinstellung beträgt 1 GB. Maximal 40 GB |
Kernelstatus anzeigen
Nachdem Sie Ihren Spark-Kernel erstellt haben, können Sie die Kerneldetails anzeigen.
Geben Sie Folgendes ein, um den Kernelstatus anzuzeigen:
curl -k -X GET <KERNEL_ENDPOINT>/<kernel_id> -H "Authorization: Bearer ${TOKEN}"
Beispielantwort:
{
"id": "<kernel_id>",
"name": "scala",
"connections": 0,
"last_activity": "2021-07-16T11:06:26.266275Z",
"execution_state": "starting"
}
Kernel löschen
Sie können einen Spark-Kernel löschen, indem Sie Folgendes eingeben:
curl -k -X DELETE <KERNEL_ENDPOINT>/<kernel_id> -H "Authorization: Bearer ${TOKEN}"
Kernel auflisten
Sie können alle aktiven Spark-Kernel auflisten, indem Sie Folgendes eingeben:
curl -k -X GET <KERNEL_ENDPOINT> -H "Authorization: Bearer ${TOKEN}"
Spark-Kernel verwenden
Sie können die Kernel-API verwenden, die von Analytics Engine Powered by Apache Sparkbereitgestellt wird. Zum Erstellen eines Kernels überprüfen Sie den Status eines Kernels und löschen einen Kernel mit einer WebSocket-Verbindung.
Der folgende Beispielcode unter Python nutzt Tornado-Bibliotheken, um Aufrufe an HTTP und WebSocket an einen Jupyter-Kernel-Gateway-Dienst zu senden. Um diesen Beispielcode auszuführen, benötigen Sie eine Laufzeitumgebung für „ Python “ mit installiertem Tornado-Paket.
So erstellen Sie eine Spark-Anwendung:
Installiere pip. Falls Sie es noch nicht installiert haben, installieren Sie das „ Python “ Tornado-Paket mit folgendem Befehl:
yum install python3-pip -y ; pip3 install tornadoErstellen Sie in einem Arbeitsverzeichnis eine Datei mit dem Namen
client.py, die den folgenden Code enthält:Der folgende Beispielcode erstellt einen Spark-Kernel vom Typ „ Scala “ und übergibt ihm einen „ Scala “-Code zur Ausführung:
from uuid import uuid4 import requests import json from tornado import gen from tornado.escape import json_encode, json_decode, url_escape from tornado.httpclient import AsyncHTTPClient, HTTPRequest from tornado.ioloop import IOLoop from tornado.websocket import websocket_connect @gen.coroutine def main(): url = "https://cpd-cpd-instance.apps.ocp-270005she9-tcvv.cloud.techzone.ibm.com/icp4d-api/v1/authorize" payload = json.dumps({ "username": "magonzalez", "password": "magonzalez" }) headers = { 'cache-control': 'no-cache', 'Content-Type': 'application/json' } response = requests.request("POST", url, headers=headers, data=payload) response_json = json.loads(response.text) token = response_json['token'] kg_http_url = "https://cpd-cpd-instance.apps.ocp-270005she9-tcvv.cloud.techzone.ibm.com/v4/analytics_engines/c74aa35e-9fcb-4f9a-8e41-59fbdfe1054a/jkg/api/kernels" kg_ws_url = "wss://cpd-cpd-instance.apps.ocp-270005she9-tcvv.cloud.techzone.ibm.com/v4/analytics_engines/c74aa35e-9fcb-4f9a-8e41-59fbdfe1054a/jkg/api/kernels" headers = {"Authorization": 'Bearer {}'.format(token) , "Content-Type": "application/json"} validate_cert = False kernel_name="scala" kernel_payload={"name":"scala","kernel_size":{"cpu":3}} print(kernel_payload) code = """ print(s"Spark Version: ${sc.version}") print(s"Application Name: ${sc.appName}") print(s"Application ID: ${sc.applicationId}") import org.apache.spark.sql.SQLContext; val sqlContext = new SQLContext(sc); val data_df_0 = sqlContext.read.format("csv").option("header", "true").option("inferSchema", "true").option("mode", "DROPMALFORMED").csv("/opt/ibm/spark/examples/src/main/resources/people.csv"); data_df_0.show(5) """ print("Using kernel gateway URL: {}".format(kg_http_url)) print("Using kernel websocket URL: {}".format(kg_ws_url)) # Remove "/" if exists in JKG url's if kg_http_url.endswith("/"): kg_http_url=kg_http_url.rstrip('/') if kg_ws_url.endswith("/"): kg_ws_url=kg_ws_url.rstrip('/') client = AsyncHTTPClient() # Create kernel # POST /api/kernels print("Creating kernel {}...".format(kernel_name)) response = yield client.fetch( kg_http_url, method='POST', headers = headers, validate_cert=validate_cert, body=json_encode(kernel_payload), connect_timeout=240, request_timeout=240 ) kernel = json_decode(response.body) kernel_id = kernel['id'] print("Created kernel {0}.".format(kernel_id)) # Connect to kernel websocket # GET /api/kernels/<kernel-id>/channels # Upgrade: websocket # Connection: Upgrade print("Connecting to kernel websocket...") ws_req = HTTPRequest(url='{}/{}/channels'.format( kg_ws_url, url_escape(kernel_id) ), headers = headers, validate_cert=validate_cert, connect_timeout=240, request_timeout=240 ) ws = yield websocket_connect(ws_req) print("Connected to kernel websocket.") # Submit code to websocket on the 'shell' channel print("Submitting code: \n{}\n".format(code)) msg_id = uuid4().hex req = json_encode({ 'header': { 'username': '', 'version': '5.0', 'session': '', 'msg_id': msg_id, 'msg_type': 'execute_request' }, 'parent_header': {}, 'channel': 'shell', 'content': { 'code': code, 'silent': False, 'store_history': False, 'user_expressions': {}, 'allow_stdin': False }, 'metadata': {}, 'buffers': {} }) # Send an execute request ws.write_message(req) print("Code submitted. Waiting for response...") # Read websocket output until kernel status for this request becomes 'idle' kernel_idle = False while not kernel_idle: msg = yield ws.read_message() msg = json_decode(msg) msg_type = msg['msg_type'] print ("Received message type: {}".format(msg_type)) if msg_type == 'error': print('ERROR') print(msg) break # evaluate messages that correspond to our request if 'msg_id' in msg['parent_header'] and \ msg['parent_header']['msg_id'] == msg_id: if msg_type == 'stream': print(" Content: {}".format(msg['content']['text'])) elif msg_type == 'status' and \ msg['content']['execution_state'] == 'idle': kernel_idle = True # close websocket ws.close() # Delete kernel # DELETE /api/kernels/<kernel-id> print("Deleting kernel...") yield client.fetch( '{}/{}'.format(kg_http_url, kernel_id), method='DELETE', headers = headers, validate_cert=validate_cert, ) print("Deleted kernel {0}.".format(kernel_id)) if __name__ == '__main__': IOLoop.current().run_sync(main)Dabei gilt:
<token>ist das Zugriffstoken, das Sie im Abschnitt Kernel-API verwendengeneriert haben.<kernel_endpoint>ist der Spark-Kernelendpunkt, den Sie in Kernel-API verwendenerhalten haben. Beispiel für einen Kernelendpunkt:https://<cp4d_route>/v4/analytics_engines/<instance_id>/jkg/api/kernels<ws_kernel_endpoint>ist der WebSocket -Endpunkt, den Sie erstellt haben, indem Sie den Spark-Kernelendpunkt als HTTPS-Präfix in wss geändert haben.
Führen Sie die Scriptdatei aus:
python client.py
Die folgenden Code-Snippets zeigen, wie der Kernelname, die Kernelnutzdaten und Codevariablen in der Datei client.py für Python -und R-Kernel geändert werden können:
Python 3.10:
kernel_name="python310" kernel_payload={"name":"python310"} print(kernel_payload) code = '\n'.join(( "print(\"Spark Version: {}\".format(sc.version))", "print(\"Application Name: {}\".format(sc._jsc.sc().appName()))", "print(\"Application ID: {} \".format(sc._jsc.sc().applicationId()))", "sc.parallelize([1,2,3,4,5]).count()" ))R 4.2:
kernel_name="r42" kernel_payload={"name":"r42"} code = """ cat("Spark Version: ", sparkR.version()) conf = sparkR.callJMethod(spark, "conf") cat("Application Name: ", sparkR.callJMethod(conf, "get", "spark.app.name")) cat("Application ID:", sparkR.callJMethod(conf, "get", "spark.app.id")) df <- as.DataFrame(list(1,2,3,4,5)) cat(count(df)) """