运行 Spark Streaming 应用程序
Spark Streaming 支持对实时数据流 (例如来自日志文件或状态更新消息) 进行可扩展的高吞吐量容错流处理。 HDFS 目录、 TCP 套接字以及 Kafka 是Spark Streaming支持的部分数据源。 Spark 中的机器学习和图形处理方法甚至可以在数据流上使用,处理后的数据可以存储在数据库和文件中。
Spark Streaming 采用实时输入数据流,并将流分成多个批次,由 Spark 引擎进行处理,以提供最终批次的结果。 有关详细信息,请参阅 Apache Spark Streaming Programming Guide。
与 Kafka 集成
IBM Analytics Engine powered by Apache Spark 支持将 Kafka 用作实时数据流的数据源。 数据会在流式传输时进行处理,并可存储在 HDFS、数据库或仪表盘中。
kafka-stream-example.py本节将向您展示如何在名为 的示例应用程序中,结合 Kafka 在 IBM Analytics Engine powered by Apache Spark 上利用 Spark Streaming。
样本 Python Spark Streaming 应用程序:
#!/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()
运行 Spark 应用程序
要使用流经 Kafka的数据运行 Spark 应用程序 kafka-stream-example.py ,需要将所需的 Kafka 和 Spark 流式库预装入到 Spark。
IBM Analytics Engine powered by Apache Spark 提供了多种选项,可用于持久化应用程序中可能需要的任何库,包括应用程序文件。 有关详细信息,请参阅 定制服务卷实例中的 Spark 应用程序。
在以下示例中,您将需要的库和 Spark 应用程序文件上载到卷服务实例。
要运行 Spark 应用程序 kafka-stream-example.py:
从 Maven 下载以下 Python 软件包:
spark-sql-kafka; Spark V 3.3; 从 spark-sql-kafka-0-10_2.12-3.3.0.jar 下载spark-streaming-kafka; Spark V 3.3; 从 spark-streaming-kafka-0-10-assembly_2.12-3.3.0.jar 下载commons-pool2; Spark V 3.3; 从 commons-pool2-2.11.1.jar 下载
将软件包和 Spark 应用程序文件上载到卷服务实例:
通过使用 API。 有关指示信息,请参阅 定制服务卷实例中的 Spark 应用程序。
通过用户界面:
- 从 Cloud Pak for Data中的导航菜单
,单击 服务> 实例,找到服务卷实例,然后单击该实例以查看实例详细信息。
- 单击 文件浏览器 选项卡,然后上载 Spark 应用程序和您下载的 Spark Streaming JAR 文件。
- 从 Cloud Pak for Data中的导航菜单
准备 Spark 应用程序有效内容。
您需要在有效内容中定义
volumes部分,并添加卷服务实例和安装详细信息以在 Spark 应用程序启动之前装入所需的 Python 包。在以下样本
payload.json中, Spark 应用程序kafka-stream-example.py和 Kafka 库存储在安装到/myapp的data-vol卷中。 Python 应用程序和jars选项中包含的 JAR 的逗号分隔列表将自动传输到集群。{ "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": "" }\] }提交 PySpark 应用程序。 有关详细信息,请参阅 通过 API 提交 Spark 作业。