【问题标题】:Unable to receive data from Kafka to Spark streaming无法从 Kafka 接收数据到 Spark 流
【发布时间】:2021-02-21 10:32:10
【问题描述】:

我正在尝试使用 Eclipse IDE 中的 java 代码通过 kafka 生产者生成一些随机数据。我在 kafka 消费者中收到相同的数据,这些数据也是在同一个 IDE 中使用 java 代码创建的。我的工作依赖于流数据。所以,我需要火花流来接收kafka生成的随机数据。对于火花流,我在 jupyter-notebook 中使用 python 代码。要将 kafka 与 spark 集成,必须将“spark-streaming-kafka-0-10_2.12-3.0.0.jar”文件添加到 spark jar 中。我还尝试在 pyspark 中添加 jar 文件。这是我的火花代码

import time
from pyspark import SparkContext, SparkConf
from pyspark.sql import SparkSession
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
n_secs = 3
topic = "generate"
spark = SparkSession.builder.master("local[*]") \
        .appName("kafkaStreaming") \
        .config("/home/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/pyspark/spark-streaming-kafka-0-10_2.12-3.0.0.jar") \
        .getOrCreate()
sc = spark.sparkContext

ssc = StreamingContext(sc, n_secs)
kStream = KafkaUtils.createDirectStream(ssc, [topic], {
                        'bootstrap.servers':'localhost:9092',
                        'group.id':'test-group',
                        'auto.offset.reset':'latest'})
lines = kStream.map(lambda x: x[1])
words = lines.flatmap(lambda line: line.split(" "))
print(words)
ssc.start()
time.sleep(100)
ssc.stop(stopSparkContext=True,stopGraceFully=True)

在上面的代码中,我使用 SparkSession.config() 方法添加了 jar 文件。创建 DStream 后,我试图通过提供主题名称、引导服务器等使用 KafkaUtils.createDirectStream() 从 kafka 接收数据。之后,我将数据转换为 rdd 并打印结果。这是我工作的整体流程。起初,我在 java 中执行 kafka 生产者代码,它会生成一些数据并由 kafka 消费者消费。到目前为止,它工作正常。在 python 中执行 spark 流代码时,它显示了一些类似这样的错误

ERROR:root:Exception while sending command.
Traceback (most recent call last):
  File "/home/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1159, in send_command
    raise Py4JNetworkError("Answer from Java side is empty")
py4j.protocol.Py4JNetworkError: Answer from Java side is empty

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/home/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 985, in send_command
    response = connection.send_command(command)
  File "/home/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1164, in send_command
    "Error while receiving", e, proto.ERROR_ON_RECEIVE)
py4j.protocol.Py4JNetworkError: Error while receiving

Py4JError                                 Traceback (most recent call last)
<ipython-input-17-873ece723182> in <module>
     36                         'bootstrap.servers':'localhost:9092',
     37                         'group.id':'test-group',
---> 38                         'auto.offset.reset':'latest'})
     39 
     40 lines = kStream.map(lambda x: x[1])

~/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/pyspark/streaming/kafka.py in createDirectStream(ssc, topics, kafkaParams, fromOffsets, keyDecoder, valueDecoder, messageHandler)
    144             func = funcWithoutMessageHandler
    145             jstream = helper.createDirectStreamWithoutMessageHandler(
--> 146                 ssc._jssc, kafkaParams, set(topics), jfromOffsets)
    147         else:
    148             ser = AutoBatchedSerializer(PickleSerializer())

~/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py in __call__(self, *args)
   1255         answer = self.gateway_client.send_command(command)
   1256         return_value = get_return_value(
-> 1257             answer, self.gateway_client, self.target_id, self.name)
   1258 
   1259         for temp_arg in temp_args:

~/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/pyspark/sql/utils.py in deco(*a, **kw)
     61     def deco(*a, **kw):
     62         try:
---> 63             return f(*a, **kw)
     64         except py4j.protocol.Py4JJavaError as e:
     65             s = e.java_exception.toString()

~/Downloads/Spark/spark-2.4.6-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name)
    334             raise Py4JError(
    335                 "An error occurred while calling {0}{1}{2}".
--> 336                 format(target_id, ".", name))
    337     else:
    338         type = answer[1]

Py4JError: An error occurred while calling o270.createDirectStreamWithoutMessageHandler

请任何人帮助我摆脱这个问题...

【问题讨论】:

    标签: python apache-spark pyspark apache-kafka spark-streaming


    【解决方案1】:

    我可以从代码本身看到一些东西:

    • 您的 jar 工件适用于 spark 3.0,并且您使用的是 spark 版本 2.4.6。(提示:文件名中的最后 3 位数字是 spark 版本)
    • 您已在配置选项中添加了 jar 文件。我建议首先验证您正在使用的 jar 文件是否是您需要的,方法是在 spark-submit 命令中使用它作为--jar &lt;jar-file-path&gt; 。
    • 先尝试打印您的直接流,而不是对其进行各种转换。你可以这样做:
    kStream = KafkaUtils.createDirectStream(ssc, [topic], {
                            'bootstrap.servers':'localhost:9092',
                            'group.id':'test-group',
                            'auto.offset.reset':'latest'})
    kStream.pprint()
    ssc.start()
    # stream will run for 50 sec
    ssc.awaitTerminationOrTimeout(50)
    ssc.stop()
    sc.stop()
    
    • 验证您正在获取数据后,您可以使用 foreachRDD、transform 或其他 api 来处理您的数据

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-02-05
      • 1970-01-01
      • 2017-01-25
      • 2016-08-18
      • 2018-08-20
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多