【问题标题】:Sentence similarity with SparkNLP only works on Google Dataproc with ONE sentence, FAILS when multiple sentences are provided与 SparkNLP 的句子相似性仅适用于具有一个句子的 Google Dataproc,当提供多个句子时失败
【发布时间】:2023-03-16 23:10:01
【问题描述】:

将以下 colab python 代码(请参阅下面的链接)部署到 Google Cloud 上的 Dataproc,并且仅当 input_list 是一个包含一个项目的数组时,当 input_list有两个项目,然后 PySpark 作业在下面的 get_similarity 方法中的“for r in result.collect()”行中出现以下错误:

java.io.IOException: Premature EOF from inputStream
        at org.apache.hadoop.io.IOUtils.readFully(IOUtils.java:194)
        at org.apache.hadoop.hdfs.protocol.datatransfer.PacketReceiver.doReadFully(PacketReceiver.java:213)
        at org.apache.hadoop.hdfs.protocol.datatransfer.PacketReceiver.doRead(PacketReceiver.java:134)
        at org.apache.hadoop.hdfs.protocol.datatransfer.PacketReceiver.receiveNextPacket(PacketReceiver.java:109)
        at org.apache.hadoop.hdfs.server.datanode.BlockReceiver.receivePacket(BlockReceiver.java:446)
        at org.apache.hadoop.hdfs.server.datanode.BlockReceiver.receiveBlock(BlockReceiver.java:702)
        at org.apache.hadoop.hdfs.server.datanode.DataXceiver.writeBlock(DataXceiver.java:739)
        at org.apache.hadoop.hdfs.protocol.datatransfer.Receiver.opWriteBlock(Receiver.java:124)
        at org.apache.hadoop.hdfs.protocol.datatransfer.Receiver.processOp(Receiver.java:71)
        at org.apache.hadoop.hdfs.server.datanode.DataXceiver.run(DataXceiver.java:232)
        at java.lang.Thread.run(Thread.java:745)
input_list=["no error"]                 <---- works
input_list=["this", "throws EOF error"] <---- does not work

链接到 colab 使用 spark-nlp 进行句子相似性: https://colab.research.google.com/github/JohnSnowLabs/spark-nlp-workshop/blob/master/tutorials/streamlit_notebooks/SENTENCE_SIMILARITY.ipynb#scrollTo=6E0Y5wtunFi4

def get_similarity(input_list):
    df = spark.createDataFrame(pd.DataFrame({'text': input_list}))
    result = light_pipeline.transform(df)
    embeddings = []
    for r in result.collect():
        embeddings.append(r.sentence_embeddings[0].embeddings)
    embeddings_matrix = np.array(embeddings)
    return np.matmul(embeddings_matrix, embeddings_matrix.transpose())

我尝试在 hadoop 集群配置中将 dfs.datanode.max.transfer.threads 更改为 8192,但仍然没有成功:

hadoop_config.set('dfs.datanode.max.transfer.threads', "8192")

input_list 数组中有多个项目时,如何使此代码正常工作?

【问题讨论】:

  • 您是从 HDFS 还是 GCS 读取数据?
  • 最初来自 GCS,但后来在 HDFS 中
  • 两者都失败了吗?这似乎是一些并发问题。
  • 失败出现在 HDFS 上,似乎阵列被拆分为多个 HDFS 实例,这导致了并发问题

标签: apache-spark hadoop hdfs google-cloud-dataproc johnsnowlabs-spark-nlp


【解决方案1】:

java.io.IOException: Premature EOF from inputStream 可能表示磁盘带宽不足、HDFS 数据节点过载或许多其他问题:Hadoop MapReduce job I/O Exception due to premature EOF from inputStream

在 Spark 应用程序中增加 DataNode 传输线程的数量并没有改变任何东西,因为您需要在每个集群工作人员的 HDFS 配置中更改此属性并在每个工作人员上重新启动 DataNode 服务。最简单的方法是使用 hdfs:dfs.datanode.max.transfer.threads=8192 cluster property 重新创建一个集群。

请注意,如果问题的根本原因是磁盘带宽不足,那么在 DataNodes 中增加传输线程的数量只会夸大它,而不是修复它。

您有多种选择来尝试解决此问题:

  1. 要增加本地磁盘带宽,请在创建集群时在工作节点上使用PD-SSDLocal SSD

  2. 如果您使用具有少量工作人员的集群,HDFS 数据节点(每个工作人员 1 个)可能无法处理负载,作为一种解决方法,您可以增加集群中工作人员的数量或重新创建如果您使用具有超过 4 个 CPU 内核的工作线程,则集群具有相同容量但较小数量的工作线程。

  3. 使用 Google Cloud Storage(gs:// 架构)而不是 HDFS 来存储您处理的数据 - Google Cloud Storage 的扩展性比 HDFS 好得多,并且应该开箱即用。

【讨论】:

    猜你喜欢
    • 2018-10-01
    • 2016-08-14
    • 1970-01-01
    • 2019-11-27
    • 2021-11-21
    • 2020-08-10
    • 2021-05-23
    • 2022-06-15
    • 2018-02-02
    相关资源
    最近更新 更多