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:

  1. Genera un token di accesso se non ne hai uno. Vedi Genera un token di accesso.

  2. Esporta il token in una variabile:

    export TOKEN=<token generated>
    
  3. 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.

  4. Ottieni l'endpoint kernel Spark per l'istanza.

    1. Dal menu di navigazione Menu di navigazione Cloud Pak for Data 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.
    2. In Informazioni di accesso, copiare e salvare l'endpoint kernel Spark.
  5. 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:

  1. Installare pip. Se non lo hai, installa il pacchetto Python Tornado utilizzando questo comando:

    yum install  python3-pip -y ; pip3 install tornado
    
  2. In una directory di lavoro, creare un file denominato client.py contenente 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.
  3. 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))
    """