【问题标题】:Spark-streaming application hangs when I use yarn-mode当我使用纱线模式时,Spark-streaming 应用程序挂起
【发布时间】:2017-03-06 04:11:39
【问题描述】:

我对纱线上的火花流有疑问。

当我在本地模式下启动我的脚本时,它运行良好:我可以从 Flume 接收和打印事件。

from pyspark.streaming.flume import FlumeUtils
from pyspark.streaming import StreamingContext
from pyspark.storagelevel import StorageLevel
from pyspark import SparkConf, SparkContext

conf = SparkConf().setAppName("Test").setMaster('local[*]')
sc = SparkContext(conf=conf)

ssc = StreamingContext(sc, 1)
hostname = 'myhost.com'
port = 6668
addresses = [(hostname, port)]

flumeStream = FlumeUtils.createPollingStream(ssc, addresses, \
                                         storageLevel=StorageLevel(True, True, False, False, 2), \
                                         maxBatchSize=1000, parallelism=5)

flumeStream.pprint()

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

输出:

-------------------------------------------
Time: 2017-03-03 00:49:34
-------------------------------------------

-------------------------------------------
Time: 2017-03-03 00:49:35
-------------------------------------------

17/03/03 00:49:35 WARN storage.BlockManager: Block input-0-1488476966735 replicated to only 0 peer(s) instead of 1 peers
-------------------------------------------
Time: 2017-03-03 00:49:36
-------------------------------------------
({u'timestamp': u'1488476971253', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.262000]; 2; 3341678; 3279.39; 97')
({u'timestamp': u'1488476971265', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.274000]; 4; 2690399; 69.24; 27')
({u'timestamp': u'1488476971276', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.285000]; 6; 7266957; 514.57; 25')
({u'timestamp': u'1488476971286', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.296000]; 8; 9220339; 3189.55; 5')
({u'timestamp': u'1488476971298', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.307000]; 10; 2897030; 1029.84; 56')
({u'timestamp': u'1488476971308', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.317000]; 12; 4710976; 1125.88; 35')
({u'timestamp': u'1488476971340', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.349000]; 14; 4894562; 707.43; 50')
({u'timestamp': u'1488476971371', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.380000]; 16; 7370409; 1056.91; 1')
({u'timestamp': u'1488476971402', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.411000]; 18; 6669529; 2868.7; 56')
({u'timestamp': u'1488476971433', u'Severity': u'4', u'Facility': u'3'}, u'<28>[2017-03-03 00:49:31.442000]; 20; 7823207; 791.02; 15')
...

-------------------------------------------
Time: 2017-03-03 00:49:37
-------------------------------------------

但如果我尝试以纱线模式开始,例如:

conf = SparkConf().setAppName("Test")

conf = SparkConf().setAppName("Test").set("spark.executor.memory", "1g").set("spark.driver.memory", "2g")

然后(似乎)我的应用程序进入无限循环,可能正在等待某些东西。屏幕上的图像挂起如下:

-------------------------------------------
Time: 2017-03-03 00:59:34
-------------------------------------------

Spark-tasks 没有 spark-streaming 的工作正常在纱线和本地模式下。

我的基础架构是:

  • 第一个节点:job history server和ResourceManager,在同一个地方也是flume agent。

  • 第二个和第三个节点 - 节点管理器。

我正在从节点 2 启动我的 spark-streaming 应用程序。从日志中我看到节点 3 和资源管理器之间的连接正常工作。

感谢任何想法,任何帮助!谢谢!

【问题讨论】:

    标签: apache-spark pyspark spark-streaming


    【解决方案1】:

    通过在集群中再添加一个节点解决了问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-12-13
      • 1970-01-01
      • 2020-05-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-06-02
      相关资源
      最近更新 更多