Esecuzione interattiva delle applicazioni Spark
Puoi eseguire le tue applicazioni Spark in modo interattivo utilizzando l'API Kernel.
Kernel as a service fornisce:
- Kernel Jupyter come entità di prima classe
- Un cluster dedicato per kernel
- Cluster e Kernel personalizzati con le librerie utente
Utilizzo dell'API Kernel
Puoi eseguire un'applicazione Spark in modo interattivo utilizzando l'API Kernel. Ogni applicazione viene eseguita in un kernel in un cluster dedicato. Tutte le impostazioni di configurazione che passi attraverso l'API sovrascrivono le configurazioni predefinite.
Per eseguire un'applicazione Spark in modo interattivo:
Genera un token di accesso se non ne hai uno. Vedi Genera un token di accesso.
Esporta il token in una variabile:
export TOKEN=<token generated>Configurare un'istanza di Analytics Engine powered by Apache Spark. Hai bisogno di un ruolo di Amministratore nel progetto o nello spazio di distribuzione per eseguire il provisioning di un'istanza. Vedi Provisioning di una istanza.
Ottieni l'endpoint kernel Spark per l'istanza.
- Dal menu di navigazione
in Cloud Pak for Data, fai clic su Servizi> Istanze, trova l'istanza e fai clic su di essa per visualizzare i dettagli dell'istanza.
- In Informazioni di accesso, copiare e salvare l'endpoint kernel Spark.
- Dal menu di navigazione
Crea un kernel utilizzando l'endpoint kernel e il token di accesso che hai generato. Questo esempio include i parametri minimi obbligatori necessari:
curl -k -X POST <KERNEL_ENDPOINT> -H "Authorization: Bearer ${TOKEN}" -d '{"name":"scala" }'
Crea JSON del kernel e parametri supportati
Se si utilizza il servizio Kernel per avviare l'API REST di creazione del kernel, è possibile aggiungere configurazioni avanzate al payload.
Considera il seguente payload di esempio e l'elenco di parametri che puoi passare nell'API di creazione kernel:
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"
}
}
}
}'
Risposta:
{
"id": "<kernel_id>",
"name": "scala",
"connections": 0,
"last_activity": "2021-07-16T11:06:26.266275Z",
"execution_state": "starting"
}
API kernel Spark utilizzando il runtime spark personalizzato
Un esempio di payload di input per la modifica della versione di runtime spark:
{
"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"
}
}
}
}
Parametri API kernel Spark
Questi sono i parametri che puoi utilizzare nell'API Kernel:
| Nome | Obbligatorio/Facoltativo | Tipo | Descrizione |
|---|---|---|---|
| Nome | Obbligatorio | Stringa | Specifica il nome kernel. I valori supportati sono: scala, r, r42, python39, python310 |
| motore | Facoltativo | Coppie chiave-valore | Specifica il runtime Spark con informazioni sulla configurazione e sulla versione |
| engine.runtime.spark_version | Facoltativo | Stringa | Specifica la versione di runtime Spark da utilizzare per il kernel. IBM Cloud Pak for Data supporta Spark 3.4. |
| tipo | Obbligatorio se è specificato il motore | Stringa | Specifica il tipo di runtime del kernel. Attualmente, è supportato solo spark . |
| conf | Facoltativo | Oggetto JSON valore chiave | Specifica i valori di configurazione Spark che sovrascrivono i valori predefiniti |
| env | Facoltativo | Oggetto JSON valore chiave | Specifica le variabili di ambiente Spark richieste per il lavoro |
| dimensione | Facoltativo | Prende i parametri num_workers, worker_size e master_size |
|
| num_worker | Obbligatorio se è specificata la dimensione | Numero intero | Specifica il numero di nodi di lavoro nel cluster Spark. num_workers è uguale al numero di esecutori desiderato. Il valore predefinito è 1 executor per nodo di lavoro. Il numero massimo di executor supportati è 50. |
| dimensione_lavoro | Facoltativo | Prende i parametri cpu e memory. |
|
| cpu | Obbligatorio se è specificato worker_size |
Numero intero | Specifica la quantità di CPU per il nodo di lavoro. Il valore predefinito è 1 CPU. Il massimo è 10 CPU |
| Memoria | Obbligatorio se è specificato worker_size |
Numero intero | Specifica la quantità di memoria per ogni nodo di lavoro. Il valore predefinito è 1 GB. Il valore massimo è 40 GB |
| dimensione_master | Facoltativo | Acquisisce i parametri cpu e memory |
|
| cpu | Obbligatorio se è specificato master_size |
Numero intero | Specifica la quantità di CPU per il nodo master. Il valore predefinito è 1 CPU. Il massimo è 10 CPU |
| Memoria | Obbligatorio se è specificato master_size |
Numero intero | Specifica la quantità di memoria per ciascun nodo master. Il valore predefinito è 1 GB. Il valore massimo è 40 GB |
| dimensione_kernel | Facoltativo | Acquisisce i parametri cpu e memory |
|
| cpu | Obbligatorio se è specificato kernel_size |
Numero intero | Specifica la quantità di CPU per il nodo kernel. Il valore predefinito è 1 CPU. Il massimo è 10 CPU |
| Memoria | Obbligatorio se è specificato kernel_size |
Numero intero | Specifica la quantità di memoria per ciascun nodo kernel. Il valore predefinito è 1 GB. Il valore massimo è 40 GB |
Visualizzazione dello stato del kernel
Dopo aver creato il kernel Spark, è possibile visualizzarne i dettagli.
Per visualizzare lo stato del kernel, immettere:
curl -k -X GET <KERNEL_ENDPOINT>/<kernel_id> -H "Authorization: Bearer ${TOKEN}"
Esempio di risposta:
{
"id": "<kernel_id>",
"name": "scala",
"connections": 0,
"last_activity": "2021-07-16T11:06:26.266275Z",
"execution_state": "starting"
}
Eliminazione di un kernel
È possibile eliminare un kernel Spark immettendo quanto segue:
curl -k -X DELETE <KERNEL_ENDPOINT>/<kernel_id> -H "Authorization: Bearer ${TOKEN}"
Elenco dei kernel
È possibile elencare tutti i kernel Spark attivi immettendo quanto segue:
curl -k -X GET <KERNEL_ENDPOINT> -H "Authorization: Bearer ${TOKEN}"
Utilizzo dei kernel Spark
Puoi utilizzare l'API Kernel fornita da Analytics Engine Powered by Apache Spark. Per creare un kernel, controllare lo stato di un kernel ed eliminare un kernel con una connessione websocket.
Il seguente codice di esempio Python utilizza le librerie Tornado per effettuare chiamate HTTP e WebSocket a un servizio Jupyter Kernel Gateway. Per eseguire questo codice di esempio è necessario un ambiente di runtime Python con il pacchetto Tornado installato.
Per creare un'applicazione Spark:
Installare pip. Se non lo hai, installa il pacchetto Python Tornado utilizzando questo comando:
yum install python3-pip -y ; pip3 install tornadoIn una directory di lavoro, creare un file denominato
client.pycontenente il codice seguente.Il codice di esempio riportato di seguito crea un kernel Spark di tipo " Scala " e invia ad esso un codic Scala e da eseguire:
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)dove:
<token>è il token di accesso che hai generato in Utilizzo dell'API Kernel.<kernel_endpoint>è l'endpoint del kernel Spark che hai ottenuto in Utilizzo dell'API Kernel. Esempio di endpoint kernel:https://<cp4d_route>/v4/analytics_engines/<instance_id>/jkg/api/kernels.<ws_kernel_endpoint>è l'endpoint WebSocket creato prendendo l'endpoint del kernel Spark e modificando il prefisso https in wss.
Eseguire il file di script:
python client.py
Ecco i frammenti di codice che mostrano come il nome kernel, il payload del kernel e le variabili di codice possono essere modificati nel file client.py per i kernel Python e R:
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)) """