好久没做Spark了,让我来帮你吧!
首先,当您使用 Spark 3.0.0 时,您可以使用 Spark Structured Streaming,该 API 基于数据帧将更易于使用。如您所见here in the link of the docs,有一个 kafka 与 PySpark 在结构化流模式下的集成指南。
就像这个查询一样简单:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "test_topic") \
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
然后,您可以使用 ML 管道来使用此数据帧,以应用您需要的一些 ML 技术和模型。正如您在this DataBricks notebook 中看到的那样,他们提供了一些使用 ML 进行结构化流式传输的示例。这是用 Scala 编写的,但它将是一个很好的灵感来源。你可以把它和ML PySpark docs结合起来,用Python翻译它
编辑:为了使其在 PySpark 和 Kafka 之间工作而应遵循的实际步骤
1 - 卡夫卡设置
所以首先我设置了我的本地 Kafka:
wget https://archive.apache.org/dist/kafka/0.10.2.2/kafka_2.12-0.10.2.2.tgz
tar -xzf kafka_2.11-0.10.2.0.tgz
我打开 4 个 shell,运行 zookeeper/server/create_topic/write_topic 脚本:
cd kafka_2.11-0.10.2.0
bin/zookeeper-server-start.sh config/zookeeper.properties
cd kafka_2.11-0.10.2.0
bin/kafka-server-start.sh config/server.properties
cd kafka_2.11-0.10.2.0
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic test
bin/kafka-topics.sh --list --zookeeper localhost:2181
- 主题中的测试消息(以交互方式将它们写入 shell 以进行测试):
cd kafka_2.11-0.10.2.0
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test
2 - PySpark 设置
获取额外的 jars
现在我们已经设置了 Kafka,我们将使用特定的 jars 下载设置 PySpark:
- spark-streaming-kafka-0-10-assembly_2.12-3.0.0.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-streaming-kafka-0-10-assembly_2.12/3.0.0/spark-streaming-kafka-0-10-assembly_2.12-3.0.0.jar
- spark-sql-kafka-0-10_2.12-3.0.0.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_2.12/3.0.0/spark-sql-kafka-0-10_2.12-3.0.0.jar
wget https://repo1.maven.org/maven2/org/apache/commons/commons-pool2/2.8.0/commons-pool2-2.8.0.jar
- kafka-clients-0.10.2.2.jar
wget https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/0.10.2.2/kafka-clients-0.10.2.2.jar
运行 PySpark shell 命令
如果执行pyspark命令时不在jars文件夹中,不要忘记为每个jars指定文件夹路径。
PYSPARK_PYTHON=python3 $SPARK_HOME/bin/pyspark --jars spark-sql-kafka-0-10_2.12-3.0.0.jar,spark-streaming-kafka-0-10-assembly_2.12-3.0.0.jar,kafka-clients-0.10.2.2.jar,commons-pool2-2.8.0.jar
3 - 运行 PySpark 代码
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "test") \
.load()
query = df \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
.writeStream \
.format("console") \
.start()
干杯