【问题标题】:How to load Kafka topic data into a Spark Dstream in Python如何在 Python 中将 Kafka 主题数据加载到 Spark Dstream 中
【发布时间】:2020-11-26 10:25:39
【问题描述】:

我正在使用带有 Python 的 Spark 3.0.0。 我在 Kafka 中有一个 test_topic,它正在从 csv 生成。

下面的代码从该主题消费到 Spark,但我在某处读到它需要在 DStream 中,然后我才能对其进行任何 ML。

import json
from json import loads
from kafka import KafkaConsumer
from pyspark import SparkContext
from pyspark.streaming import StreamingContext

sc = SparkContext("local[2]", "test")
ssc = StreamingContext(sc, 1)

consumer = KafkaConsumer('test_topic',
                    bootstrap_servers =['localhost:9092'],
                    api_version=(0, 10))

消费者返回<kafka.consumer.group.KafkaConsumer at 0x13bf55b0>

如何编辑上面的代码给我一个 DStream?

我是新手,所以请指出任何愚蠢的错误。

编辑: 以下是我的生产者代码:

import json
import csv
from json import dumps
from kafka import KafkaProducer
from time import sleep

producer = KafkaProducer(bootstrap_servers=['localhost:9092'])
value_serializer=lambda x:dumps(x)

with open('test_data.csv') as file:
reader = csv.DictReader(file, delimiter=';')
for row in reader:
    producer.send('test_topic', json.dumps(row).encode('utf=8'))
    sleep(2)
    print ('Message sent ', row)

【问题讨论】:

    标签: apache-spark pyspark apache-kafka


    【解决方案1】:

    好久没做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
    
    • commons-pool2-2.8.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()
    

    干杯

    【讨论】:

    • 感谢文档的链接。我遇到了以下错误 AnalysisException: Failed to find data source: kafka.请按照“结构化流+Kafka集成指南”的部署部分部署应用程序。我已经在编辑中发布了我的生产者代码。
    • 这对我有用。非常感谢。 @tricky google 用于 maven 并使用命令行“--jar” /Users/abhishek_kumar/Documents/work/code/pyspark-3.0.0/deps/bin/spark-submit --jars /Users 将其添加到 jars 中/abhishek_kumar/Documents/work/code/pyspark_jars/commons-pool2-2.8.0.jar,/Users/abhishek_kumar/Documents/work/code/pyspark_jars/spark-sql-kafka-0-10_2.12-3.0.0。 jar,/Users/abhishek_kumar/Documents/work/code/pyspark_jars/kafka-clients-2.6.0.jar,/Users/abhishek_kumar/Documents/work/code/pyspark_jars/spark-token-provider-kafka-0-10_2。 12-3.0.0.jar kafka_test_1.py
    • 你能改变你的例子,在 python 中用流中的数据做一些实际的事情吗?
    【解决方案2】:

    您需要使用 org.apache.spark:spark-sql-kafka-0-10_2.12:3.0.0 包来运行它。它将使用 spark-submit 下载相关的 jars。

    【讨论】:

      【解决方案3】:

      您需要使用 KafkaUtils createDirectStream 方法。

      这是来自official Spark documentation的代码示例:

      from pyspark.streaming.kafka import KafkaUtils
       directKafkaStream = KafkaUtils.createDirectStream(ssc, [topic], {"metadata.broker.list": brokers})
      

      【讨论】:

      • 我刚刚意识到您使用的是 Spark 3.0,我对它的经验较少,但上面似乎缺少它,documentation 说:用于从 Kafka 和如果 Spark Streaming 核心 API 中不存在 Kinesis,则必须将相应的工件 spark-streaming-xyz_2.12 添加到依赖项中。例如,一些常见的如下。 spark-streaming-kafka-0-10_2.12
      猜你喜欢
      • 2017-09-30
      • 1970-01-01
      • 1970-01-01
      • 2018-03-13
      • 2017-09-03
      • 2018-01-10
      • 1970-01-01
      • 2017-08-08
      • 2023-03-29
      相关资源
      最近更新 更多