【问题标题】:Loading large sparknlp pipeline into Apache Spark batch job taking too long将大型 sparknlp 管道加载到 Apache Spark 批处理作业中花费的时间太长
【发布时间】:2021-08-24 10:37:49
【问题描述】:

我正在使用 johnsnowlabs 的 SparkNLP 从我的文本数据中提取嵌入,下面是管道。模型保存到hdfs后大小为1.8g

embeddings = BertSentenceEmbeddings.pretrained("labse", "xx") \
      .setInputCols("sentence") \
      .setOutputCol("sentence_embeddings")
nlp_pipeline = Pipeline(stages=[document_assembler, sentence_detector, embeddings])
pipeline_model = nlp_pipeline.fit(spark.createDataFrame([[""]]).toDF("text"))

我使用pipeline_model.save("hdfs:///<path>")pipeline_model 保存到HDFS

上面只执行了一次

在另一个脚本中,我使用pipeline_model = PretrainedPipeline.from_disk("hdfs:///<path>")HDFS 加载存储的管道。

上面的代码加载了模型但是占用了太多。我在 spark 本地模型(无集群)上对其进行了测试,但我拥有 94g RAM、32 个内核的高资源。

后来,我在 yarn 上部署了脚本,有 12 个 Executor,每个 Executor 有 3 个内核和 7g ram。我分配了 10g 的驱动程序内存。

脚本再次花费太多时间从 HDFS 加载保存的模型。

当火花到达这一点时(见上图),需要太多时间

我想到了一个办法

预加载

我认为的方法是以某种方式将模型预加载到内存中,当脚本想要对数据帧应用转换时,我可以以某种方式调用对预训练管道的引用并在旅途中使用它,而无需执行任何磁盘 i/o。我搜索了,但我没有找到任何地方。

请让我知道您对此解决方案的看法以及实现此目标的最佳方式。

YARN 资源

NodeName Count RAM (each) Cores (each)
Master Node 1 38g 8
Secondary Node 1 38 g 8
Worker Nodes 4 24 g 4
Total 6 172g 32

谢谢

【问题讨论】:

  • 我在 Hadoop cpu 集群上使用 sparknlp labse 时也遇到了极差的性能。最终使用了 huggingface pytorch 端口,速度提高了 X100 倍。
  • 另外,请确保您使用的是 kryo 序列化。
  • 当然 :) 与 pytorch 我只是使用 df.rdd.mapPartitions 并手动使用模型...如果您仍想使用 sparknlp,您可能需要检查 github 上的 issue #2846,关于输出不相等原模型
  • 好的,谢谢。我用了拥抱脸变压器。感谢您的提示,现在可以正常使用了
  • 请参阅示例作为答案

标签: apache-spark hadoop hadoop-yarn johnsnowlabs-spark-nlp


【解决方案1】:

正如 cmets 中所讨论的,这是基于 PyTorch 而非 SparkNLP 的解决方案。简化代码:

# labse_spark.py

LABSE_MODEL, LABSE_TOKENIZER = None


def transform(spark, df, input_col='text', output_col='output'):
    spark.sparkContext.addFile('hdfs:///path/to/labse_model')
    output_schema = T.StructType(df.schema.fields + [T.StructField(output_col, T.ArrayType(T.FloatType()))])

    rdd = df.rdd.mapPartitions(_map_partitions_func(input_col, output_col))
    res = spark.createDataFrame(data=rdd, schema=output_schema)
    return res


def _map_partitions_func(input_col, output_col):
    def executor_func(rows):
        # load everything to memory (partitions should be small, ~1k rows per partition):
        pandas_df = pd.DataFrame([r.asDict() for r in rows])
        global LABSE_MODEL, LABSE_TOKENIZER
        if not (LABSE_TOKENIZER or LABSE_MODEL):  # should happen once per executor core
            LABSE_TOKENIZER = AutoTokenizer.from_pretrained(SparkFiles.get('labse_model'))
            LABSE_MODEL = AutoModel.from_pretrained(SparkFiles.get('labse_model'))
        
        # copied from HF model card:
        encoded_input = LABSE_TOKENIZER(
            pandas_df[input_col].tolist(), padding=True, truncation=True, max_length=64, return_tensors='pt')
        with torch.no_grad():
            model_output = LABSE_MODEL(**encoded_input)
        embeddings = model_output.pooler_output
        embeddings = torch.nn.functional.normalize(embeddings)

        pandas_df[output_col] = pd.Series(embeddings.tolist())
        return pandas_df.to_dict('records')

    return executor_func

【讨论】:

    猜你喜欢
    • 2017-04-24
    • 2018-03-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-04
    • 1970-01-01
    • 2019-02-01
    相关资源
    最近更新 更多