Ejecutar aplicaciones Spark Streaming

Spark Streaming permite un procesamiento de secuencias escalable, de alto rendimiento y tolerante a errores de secuencias de datos en directo, por ejemplo desde archivos de registro o mensajes de actualización de estado. HDFS Los directorios, los sockets TCP y Kafka son algunas de las fuentes de datos compatibles con Spark Streaming. Los métodos de procesamiento de gráficos y aprendizaje automático en Spark se pueden incluso utilizar en flujos de datos y los datos procesados se pueden almacenar en bases de datos y archivos.

Spark Streaming toma secuencias de datos de entrada en directo y divide las secuencias en lotes, que son procesadas por el motor de Spark para proporcionar un lote final de resultados. Para obtener detalles, consulte la publicación Apache Spark Streaming Programming Guide.

Integración con Kafka

IBM Analytics Engine powered by Apache Spark admite Kafka como fuente de datos para la transmisión de datos en tiempo real. Los datos se procesan a medida que se transmiten y pueden almacenarse en un HDFS s, bases de datos o paneles de control.

En esta sección se explica cómo puedes aprovechar Spark Streaming en IBM Analytics Engine powered by Apache Spark junto con Kafka en una aplicación de ejemplo llamada kafka-stream-example.py.

Aplicación Spark Streaming Python de ejemplo:

#!/usr/bin/env python
# coding: utf-8

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.functions import*
from pyspark.sql.types import*
import time

#create spark session
spark = SparkSession.builder.getOrCreate()

# Connect to kafka server and read data stream
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "CHANGEME_KAFKA_SERVER") \
.option("kafka.sasl.jaas.config","org.apache.kafka.common.security.plain.PlainLoginModule required username='CHANGEME_USERNAME' password='CHANGEME_PASSWORD';") \
.option("kafka.security.protocol", "SASL_SSL") \
.option("kafka.sasl.mechanism", "PLAIN") \
.option("kafka.ssl.protocol", "TLSv1.2") \
.option("kafka.ssl.enabled.protocols", "TLSv1.2") \
.option("kafka.ssl.endpoint.identification.algorithm", "HTTPS") \
.option("subscribe", "CHANGEME_TOPIC") \
.load() \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")

df.printSchema()

# Write the input data to memory
query = df.writeStream.outputMode("append").format("memory").queryName("testk2s").option("partition.assignment.strategy", "range").start()

query.awaitTermination(30)

query.stop()

query.status

# Query data
test_result=spark.sql("select * from testk2s")
test_result.show(5)

spark.sql("select count(*) from testk2s").show()
test_result_small = spark.sql("select * from testk2s limit 5")
test_result_small.show()

Ejecución de las aplicaciones Spark

Para ejecutar la aplicación Spark kafka-stream-example.py utilizando datos que se transmiten a través de Kafka, debe precargar las bibliotecas de modalidad continua de Kafka y Spark necesarias en Spark.

IBM Analytics Engine powered by Apache Spark ofrece varias opciones para almacenar cualquier biblioteca que puedas necesitar en tus aplicaciones, incluido el archivo de la aplicación. Para obtener detalles, consulte Personalización de aplicaciones Spark en una instancia de volumen de servicio.

En el ejemplo siguiente, subirá las bibliotecas necesarias y el archivo de aplicación Spark a una instancia de servicio de volumen.

Para ejecutar la aplicación Spark kafka-stream-example.py:

  1. Descargue los siguientes paquetes Python de Maven:

  2. Cargue los paquetes y el archivo de aplicación Spark en una instancia de servicio de volumen:

    • Utilizando la API. Para obtener instrucciones, consulte Personalización de aplicaciones Spark en una instancia de volumen de servicio.

    • A través de la interfaz de usuario:

      1. En el menú de navegación Menú de navegación de Cloud Pak for Data en Cloud Pak for Data, pulse Servicios > Instancias, busque la instancia de volumen de servicio y pulse en ella para ver los detalles de la instancia.
      2. Pulse el separador Navegador de archivos y cargue la aplicación Spark y los archivos JAR de Spark Streaming que ha descargado.
  3. Prepare la carga útil de la aplicación Spark.

    Debe definir la sección volumes en la carga útil y añadir la instancia de servicio de volumen y los detalles de montaje para cargar los paquetes Python necesarios antes de que se inicie la aplicación Spark.

    En el ejemplo siguiente payload.json, la aplicación Spark kafka-stream-example.py y las bibliotecas Kafka se almacenan en el volumen data-vol que se monta en /myapp. La aplicación Python y la lista separada por comas de JAR incluidos en la opción jars se transfieren automáticamente al clúster.

    {
        "application_details": { 
            "application": "/myapp/kafka-stream-example.py", 
            "arguments": [""],
            "conf": { 
                "spark.app.name": "SparkStreams",
                "spark.eventLog.enabled": "true" 
            },
            "jars": "/myapp/spark-sql-kafka-0-10_2.12-3.3.0.jar,/myapp/spark-streaming-kafka-0-10-assembly_2.12-3.3.0.jar,/myapp/commons-pool2-2.11.1.jar" 
        }, 
        "volumes": \[{ 
            "name": "data-vol", 
            "mount_path": "/myapp", 
            "source_sub_path": "" 
        }\]
    }
    
  4. Envíe la aplicación PySpark . Para obtener detalles, consulte Envío de trabajos Spark a través de la API.

Más información