【问题标题】:Pyspark dropping all connectionsPyspark 丢弃所有连接
【发布时间】:2020-04-30 09:07:50
【问题描述】:

我正在尝试诊断我的 spark 网格上的一些奇怪的连接问题:我看到大量断开的连接。

我正在分布式 pyspark 集群上运行类似这样的东西

spark_context.parallelize(tasks)) \
                .map(lambda kwargs: my_mapped_fn(**kwargs) \
                .reduceByKey(my_reduce_by_key) \
                .map(lambda (x,y): (x, my_final_map(x,y))) \
                .reduce(my_final_reduce)

我很确定它在 my_final_map 部分失败了,因此我怀疑闭包传输,以至于我的工作失败了。

这是我得到的错误:

java.io.IOException: Failed to connect to 10.12.9.117:38103
    at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:228)
    at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:179)
    at org.apache.spark.network.netty.NettyBlockTransferService$$anon$1.createAndStart(NettyBlockTransferService.scala:97)
    at org.apache.spark.network.shuffle.RetryingBlockFetcher.fetchAllOutstanding(RetryingBlockFetcher.java:140)
    at org.apache.spark.network.shuffle.RetryingBlockFetcher.start(RetryingBlockFetcher.java:120)
    at org.apache.spark.network.netty.NettyBlockTransferService.fetchBlocks(NettyBlockTransferService.scala:106)
    at org.apache.spark.network.BlockTransferService.fetchBlockSync(BlockTransferService.scala:92)
    at org.apache.spark.storage.BlockManager.getRemoteBytes(BlockManager.scala:579)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$1.apply$mcV$sp(TaskResultGetter.scala:82)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$1.apply(TaskResultGetter.scala:63)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$1.apply(TaskResultGetter.scala:63)
    at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1951)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3.run(TaskResultGetter.scala:62)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
    at java.lang.Thread.run(Thread.java:745)
Caused by: io.netty.channel.AbstractChannel$AnnotatedConnectException: Connection refused: 10.12.9.117:38103
    at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method)
    at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:744)
    at io.netty.channel.socket.nio.NioSocketChannel.doFinishConnect(NioSocketChannel.java:257)
    at io.netty.channel.nio.AbstractNioChannel$AbstractNioUnsafe.finishConnect(AbstractNioChannel.java:291)
    at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:640)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:575)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:489)
    at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:451)
    at io.netty.util.concurrent.SingleThreadEventExecutor$2.run(SingleThreadEventExecutor.java:140)
    at io.netty.util.concurrent.DefaultThreadFactory$DefaultRunnableDecorator.run(DefaultThreadFactory.java:144)
    ... 1 more

【问题讨论】:

  • 为什么在调用结束时执行“collect()”?
  • 请分享您得到的错误、您尝试处理的数据量、集群的大小。问题的根源可能是向驱动程序发送大量数据的“collect()”。
  • 嘿@Yaron,谢谢你的帮助!我仔细检查了,我将收集切换到减少步骤并删除局部变量。仍然看到错误的想法:(

标签: apache-spark pyspark


【解决方案1】:

如果有人觉得这很有用,实际答案与 spark 完全无关。事实上,一些节点上的 IP 地址查找被破坏了。

【讨论】:

    猜你喜欢
    • 2017-09-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-07-28
    • 2021-05-01
    • 1970-01-01
    • 2020-06-01
    • 2011-07-07
    相关资源
    最近更新 更多