【问题标题】:Why does my Spark Streaming application shut down immediately (and not process any Kafka records)?为什么我的 Spark Streaming 应用程序会立即关闭(并且不处理任何 Kafka 记录)?
【发布时间】:2017-04-10 23:02:29
【问题描述】:

我按照Spark Streaming + Kafka Integration Guide (Kafka broker version 0.8.2.1 or higher) 中描述的示例在 Python 中创建了一个 Spark 应用程序,以使用 Apache Spark 流式传输 Kafka 消息,但在我有机会发送任何消息之前它就关闭了。

这是关闭部分在输出中开始的地方。

16/11/26 17:11:06 INFO BlockManagerMaster: Registered BlockManager BlockManagerId(driver, 1********6, 58045)
16/11/26 17:11:06 INFO VerifiableProperties: Verifying properties
16/11/26 17:11:06 INFO VerifiableProperties: Property group.id is overridden to 
16/11/26 17:11:06 INFO VerifiableProperties: Property zookeeper.connect is overridden to 
16/11/26 17:11:07 INFO SparkContext: Invoking stop() from shutdown hook
16/11/26 17:11:07 INFO SparkUI: Stopped Spark web UI at http://192.168.1.16:4040
16/11/26 17:11:07 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped!
16/11/26 17:11:07 INFO MemoryStore: MemoryStore cleared
16/11/26 17:11:07 INFO BlockManager: BlockManager stopped
16/11/26 17:11:07 INFO BlockManagerMaster: BlockManagerMaster stopped
16/11/26 17:11:07 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped!
16/11/26 17:11:07 INFO SparkContext: Successfully stopped SparkContext
16/11/26 17:11:07 INFO ShutdownHookManager: Shutdown hook called
16/11/26 17:11:07 INFO ShutdownHookManager: Deleting directory /private/var/folders/yn/t3pvrk7s231_11ff2lqr4jhr0000gn/T/spark-1876feee-9b71-413e-a505-99c414aafabf/pyspark-1d97c3dd-0889-42ed-b559-d0fd473faa22
16/11/26 17:11:07 INFO ShutdownHookManager: Deleting directory /private/var/folders/yn/t3pvrk7s231_11ff2lqr4jhr0000gn/T/spark-1876feee-9b71-413e-a505-99c414aafabf

有什么方法可以让它等待还是我错过了什么?

完整代码:

from pyspark.streaming.kafka import KafkaUtils
from pyspark import SparkContext
from pyspark.streaming import StreamingContext

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

directKafkaStream = KafkaUtils.createDirectStream(ssc, ["next"], {"metadata.broker.list": "localhost:9092"})

offsetRanges = []

def storeOffsetRanges(rdd):
    global offsetRanges
    offsetRanges = rdd.offsetRanges()
    return rdd

def printOffsetRanges(rdd):
    for o in offsetRanges:
        print("Printing! %s %s %s %s" % o.topic, o.partition, o.fromOffset, o.untilOffset)

directKafkaStream\
    .transform(storeOffsetRanges)\
    .foreachRDD(printOffsetRanges)

这是运行它的命令,以防万一。

spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.0.2 producer.py

【问题讨论】:

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


    【解决方案1】:

    您还需要启动流式传输上下文。看看这个例子。 http://spark.apache.org/docs/latest/streaming-programming-guide.html#a-quick-example

    ssc.start()             # Start the computation
    ssc.awaitTermination()  # Wait for the computation to terminate
    

    【讨论】:

      【解决方案2】:

      对于 Scala,在集群模式下提交到 yarn 时,我不得不使用 awaitAnyTermination:

      query.start()
      sparkSession.streams.awaitAnyTermination()
      

      按照此处的文档Structured Streaming Guide 进行快速示例的一半。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2017-04-02
        • 2020-12-04
        • 1970-01-01
        • 2022-12-28
        • 2017-02-06
        相关资源
        最近更新 更多