Spark-nlp: When loading pretrained pipeline from HDFS in offline mode, pathname error is raised

Created on 21 May 2020  路  15Comments  路  Source: JohnSnowLabs/spark-nlp

Description


I have loaded the pretrained pipeline "recognize_entities_dl" on HDFS but, when I try to load it from spark (in Standalone Mode), I get the following error:

Exception in thread "main" java.lang.IllegalArgumentException: Pathname /C:/Users/Utente/AppData/Local/Temp/spark-85c6369e-13a0-4e3b-b6f1-c4de226ba4e5/userFiles-84d376ac-4937-491b-8e86-57ee04d5761d/storage/EMBEDDINGS_glove_100d from C:/Users/Utente/AppData/Local/Temp/spark-85c6369e-13a0-4e3b-b6f1-c4de226ba4e5/userFiles-84d376ac-4937-491b-8e86-57ee04d5761d/storage/EMBEDDINGS_glove_100d is not a valid DFS filename.
    at org.apache.hadoop.hdfs.DistributedFileSystem.getPathName(DistributedFileSystem.java:196)
    at org.apache.hadoop.hdfs.DistributedFileSystem.access$000(DistributedFileSystem.java:105)
    at org.apache.hadoop.hdfs.DistributedFileSystem$18.doCall(DistributedFileSystem.java:1118)
    at org.apache.hadoop.hdfs.DistributedFileSystem$18.doCall(DistributedFileSystem.java:1114)
    at org.apache.hadoop.fs.FileSystemLinkResolver.resolve(FileSystemLinkResolver.java:81)
    at org.apache.hadoop.hdfs.DistributedFileSystem.getFileStatus(DistributedFileSystem.java:1114)
    at org.apache.hadoop.fs.FileSystem.exists(FileSystem.java:1400)
    at com.johnsnowlabs.storage.StorageHelper$.copyIndexToLocal(StorageHelper.scala:77)
    at com.johnsnowlabs.storage.StorageHelper$.sendToCluster(StorageHelper.scala:51)
    at com.johnsnowlabs.storage.StorageHelper$.load(StorageHelper.scala:27)
    at com.johnsnowlabs.storage.HasStorageModel$$anonfun$deserializeStorage$1.apply(HasStorageModel.scala:27)
    at com.johnsnowlabs.storage.HasStorageModel$$anonfun$deserializeStorage$1.apply(HasStorageModel.scala:26)
    at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
    at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
    at com.johnsnowlabs.storage.HasStorageModel$class.deserializeStorage(HasStorageModel.scala:26)
    at com.johnsnowlabs.nlp.embeddings.WordEmbeddingsModel.deserializeStorage(WordEmbeddingsModel.scala:32)
    at com.johnsnowlabs.storage.StorageReadable$class.readStorage(StorageReadable.scala:23)
    at com.johnsnowlabs.nlp.embeddings.WordEmbeddingsModel$.readStorage(WordEmbeddingsModel.scala:156)
    at com.johnsnowlabs.storage.StorageReadable$$anonfun$1.apply(StorageReadable.scala:26)
    at com.johnsnowlabs.storage.StorageReadable$$anonfun$1.apply(StorageReadable.scala:26)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead$1.apply(ParamsAndFeaturesReadable.scala:31)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead$1.apply(ParamsAndFeaturesReadable.scala:30)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$class.com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead(ParamsAndFeaturesReadable.scala:30)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$read$1.apply(ParamsAndFeaturesReadable.scala:41)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$read$1.apply(ParamsAndFeaturesReadable.scala:41)
    at com.johnsnowlabs.nlp.FeaturesReader.load(ParamsAndFeaturesReadable.scala:19)
    at com.johnsnowlabs.nlp.FeaturesReader.load(ParamsAndFeaturesReadable.scala:8)
    at org.apache.spark.ml.util.DefaultParamsReader$.loadParamsInstance(ReadWrite.scala:652)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$$anonfun$4.apply(Pipeline.scala:274)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$$anonfun$4.apply(Pipeline.scala:272)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
    at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
    at scala.collection.TraversableLike$class.map(TraversableLike.scala:234)
    at scala.collection.mutable.ArrayOps$ofRef.map(ArrayOps.scala:186)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$.load(Pipeline.scala:272)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:348)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:342)
    at org.apache.spark.ml.util.MLReadable$class.load(ReadWrite.scala:380)
    at org.apache.spark.ml.PipelineModel$.load(Pipeline.scala:332)
    at tags_extraction.tags_extraction_eng$.main(tags_extraction_eng.scala:161)
    at tags_extraction.tags_extraction_eng.main(tags_extraction_eng.scala)

Checking the code inside the library it seems that the error is raised during the loading phase of the storage folder inside stage "3_WORD_EMBEDDINGS_MODEL_48cffc8b9a76".
(If I try to load a model created by me and saved on HDFS, everything works correctly).

Expected Behavior


Load pretrained pipeline "recognize_entities_dl" from HDFS

Current Behavior


Pathname is not a valid DFS filename

Possible Solution

  1. Change sparkSession.sparkContext.hadoopConfigurations in order to save the files on HDFS instead of locally;
  2. Download in online mode on HDFS instead of locally;

Steps to Reproduce


  1. Set sparkSession for standalone cluster in client mode
  2. PipelineModel.load("hdfs://ip/recognize_entities_dl_path")

Context

Your Environment

  • Spark Version: 2.4.5
  • Spark NLP version: 2.4.5
  • Operating System and version: Ubuntu 18.04
Requires more input question

All 15 comments

Could you please paste the entire code that results in this error? Also, could you please describe the environment in detail, it doesn't seem a straightforward Hadoop environment like Cloudera/Hortonworks. We have tested everything in those and other clusters, so anything that can help to reproduce this is helpful.

import com.johnsnowlabs.nlp.pretrained.PretrainedPipeline
import com.johnsnowlabs.util.{ConfigLoader, PipelineModels}
import org.apache.hadoop.fs.FileSystem
import org.apache.spark.SparkContext
import org.apache.spark.ml.PipelineModel
import org.apache.spark.sql.{DataFrame, SparkSession}
import utils.Schemas.feedSchema

object tags_extraction_eng {

  System.setProperty("HADOOP_USER_NAME", "root")

  val (spark, sc) = start_spark_session()

  val fs = FileSystem.get(sc.hadoopConfiguration)
  println("HDFS_IMPL:", sc.hadoopConfiguration.get("fs.hdfs.impl"))  //null
  println("DEFAULT_FS:", sc.hadoopConfiguration.get("fs.defaultFS"))  // file:///
  println("FS_URI:", fs.getUri.toString)  // file:///
  println("FS_SCHEME:", fs.getScheme)  // file


  def start_spark_session(): (SparkSession, SparkContext) = {
    val spark = SparkSession.builder()
      .appName("feeds_tags_extraction_nlp_eng")
      //            .master("local[*]")
      .master("spark://remote_ip:7077")
      .config("spark.driver.host", "host_ip")
      .config("spark.deploy.mode", "client")

      .config("spark.jars", 
        "https://repo1.maven.org/maven2/com/johnsnowlabs/nlp/spark-nlp_2.11/2.4.5/spark-nlp_2.11-2.4.5.jar," +       
        "hdfs://remote_ip:remote_port/spark_project.jar"
      )

    val sc = spark.sparkContext
    sc.getConf.getAll.foreach(println)

    (spark, sc)
  }


  def main(args: Array[String]) {
    //    val pipeline = new PretrainedPipeline(downloadName = "recognize_entities_dl", lang = "en")
    val pipeline = PipelineModel.load("hdfs://remote_ip:remote_port/recognize_entities_dl_en_2.4.3_2.4_1584626752821/")
    print(pipeline)

    spark.close()
  }

}

The environment is composed of:

  • spark cluster consisting of 1 master and 4 workers;
  • hdfs consisting of 1 namenode and 1 datanode

All running on dockers' containers (on a remote server).

I have tried to run spark on my laptop with master = local [*] and everything works. On the other hand, when I set the master = "spark: // remote_ip_address"and try to load it in the cluster I get the error shown above. I think the problem may depend on the storage path that is set in "StorageHelper.scala".

I also tried to download the pipeline using PretrainedPipeline (" recognize_entities_dl ", lang =" en ") but it is downloaded on my pc instead that on spark-worker container.

@FedericoF93 The environment is not a Spark Cluster on Hadoop coupled with HDFS. In those environments the entire cluster knows all the configs. When you tried to see what is the default file system it gets back file:/// instead of hdfs:///. This is quite OK, however, you are loading the pipeline from HDFS that for some reason forces the pipeline to also extract the embeddings on the same file system which in this case it assumes HDFS and it crashes with Pathname exception. If you load it from a local file system this probably wouldn't happen. (needs to be tested)

@rohit-nlp Is there a way to pass a spark config or environment variable/application.conf (whichever is easier) to set the path/home for embeddings to be a local path to be tested? (this setup is like S3 or shared file system when the file system is not really part of the Hadoop or Spark master/slave rather than just being available/accessible to it)

@FedericoF93 you are running spark in client mode so your driver would be running on your laptop even if you provide master as remote spark cluster.
PretrainedPipeline runs in driver so depending on where your driver is running it will download the file there. you can change the default location of files though by settings "sparknlp.settings.pretrained.cache_folder" property if you set this to hdfs file path the downloaded files in PretrainedPipeline would go that folder.

few things if you could try and let us know we can drill down on the issue.
1) Is "recognize_entities_dl" the only pipeline which do not work in your setup or all spark-nlp pipelines do not work?
2) try to run in cluster mode and see if the pipeline is saved in your hdfs home and loaded from there.

  1. Is "recognize_entities_dl" the only pipeline which do not work in your setup or all spark-nlp pipelines do not work?

No, I've also tried to load the pretrained pipeline "explain_document_ml" with PipelineModel.load("hdfs://remote_ip:emote_port/explain_document_ml_en_2.4.0_2.4_1580252705962/") but I get a different error::

[ERROR] Exception while beginning fetch of 1 outstanding blocks 
java.io.IOException: Failed to connect to /10.128.16.2:38723
    at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:245)
    at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:187)
    at org.apache.spark.network.netty.NettyBlockTransferService$$anon$2.createAndStart(NettyBlockTransferService.scala:114)
    at org.apache.spark.network.shuffle.RetryingBlockFetcher.fetchAllOutstanding(RetryingBlockFetcher.java:141)
    at org.apache.spark.network.shuffle.RetryingBlockFetcher.start(RetryingBlockFetcher.java:121)
    at org.apache.spark.network.netty.NettyBlockTransferService.fetchBlocks(NettyBlockTransferService.scala:124)
    at org.apache.spark.network.BlockTransferService.fetchBlockSync(BlockTransferService.scala:98)
    at org.apache.spark.storage.BlockManager.getRemoteBytes(BlockManager.scala:757)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3$$anonfun$run$1.apply$mcV$sp(TaskResultGetter.scala:88)
    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:1945)
    at org.apache.spark.scheduler.TaskResultGetter$$anon$3.run(TaskResultGetter.scala:62)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
Caused by: io.netty.channel.AbstractChannel$AnnotatedConnectException: Connection timed out: no further information: /10.128.16.2:38723
Caused by: java.net.ConnectException: Connection timed out: no further information
    at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method)
    at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:717)
    at io.netty.channel.socket.nio.NioSocketChannel.doFinishConnect(NioSocketChannel.java:327)
    at io.netty.channel.nio.AbstractNioChannel$AbstractNioUnsafe.finishConnect(AbstractNioChannel.java:334)
    at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:688)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:635)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:552)
    at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:514)
    at io.netty.util.concurrent.SingleThreadEventExecutor$6.run(SingleThreadEventExecutor.java:1044)
    at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
    at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
    at java.lang.Thread.run(Thread.java:748)
  1. I modified the code as you suggested, but I didn't notice any difference. It seems like "spark.submit.deployMode" is not taken into account:
import com.johnsnowlabs.nlp.pretrained.PretrainedPipeline
import com.johnsnowlabs.util.{ConfigLoader, PipelineModels}
import org.apache.hadoop.fs.FileSystem
import org.apache.spark.SparkContext
import org.apache.spark.ml.PipelineModel
import org.apache.spark.sql.{DataFrame, SparkSession}
import utils.Schemas.feedSchema

object tags_extraction_eng {

  System.setProperty("HADOOP_USER_NAME", "root")

  val (spark, sc) = start_spark_session()

  val fs = FileSystem.get(sc.hadoopConfiguration)
  println("HDFS_IMPL:", sc.hadoopConfiguration.get("fs.hdfs.impl"))  //null
  println("DEFAULT_FS:", sc.hadoopConfiguration.get("fs.defaultFS"))  // file:///
  println("FS_URI:", fs.getUri.toString)  // file:///
  println("FS_SCHEME:", fs.getScheme)  // file


  def start_spark_session(): (SparkSession, SparkContext) = {
    val spark = SparkSession.builder()
      .appName("feeds_tags_extraction_nlp_eng")
      //            .master("local[*]")
      .master("spark://remote_ip:7077")
      .config("spark.driver.host", "host_ip")

      .config("spark.submit.deployMode", "cluster")
      .config("sparknlp.settings.pretrained.cache_folder", "hdfs://remote_ip:remote_port/cache_folder")

      .config("spark.jars", 
        "https://repo1.maven.org/maven2/com/johnsnowlabs/nlp/spark-nlp_2.11/2.4.5/spark-nlp_2.11-2.4.5.jar," +       
        "hdfs://remote_ip:remote_port/spark_project.jar"
      )

    val sc = spark.sparkContext
    sc.getConf.getAll.foreach(println)

    (spark, sc)
  }


  def main(args: Array[String]) {
    //    val pipeline = new PretrainedPipeline(downloadName = "recognize_entities_dl", lang = "en")
    val pipeline = PipelineModel.load("hdfs://remote_ip:remote_port/recognize_entities_dl_en_2.4.3_2.4_1584626752821/")
    print(pipeline)

    spark.close()
  }

}

Moreover, I tried to use new PretrainedPipeline(downloadName = "recognize_entities_dl", lang = "en") both in client and cluster mode and it returns the same error message:

explain_document_ml download started this may take some time.
Approximate size to download 9,4 MB
Download done! Loading the resource.

[Stage 0:>                                                          (0 + 1) / 1]
[ WARN] Lost task 0.0 in stage 0.0 (TID 0, 10.128.16.2, executor 0): java.io.FileNotFoundException: File file:/C:/Users/Utente/cache_pretrained/explain_document_ml_en_2.4.0_2.4_1580252705962/metadata/part-00000 does not exist
    at org.apache.hadoop.fs.RawLocalFileSystem.deprecatedGetFileStatus(RawLocalFileSystem.java:611)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileLinkStatusInternal(RawLocalFileSystem.java:824)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileStatus(RawLocalFileSystem.java:601)
    at org.apache.hadoop.fs.FilterFileSystem.getFileStatus(FilterFileSystem.java:421)
    at org.apache.hadoop.fs.ChecksumFileSystem$ChecksumFSInputChecker.<init>(ChecksumFileSystem.java:142)
    at org.apache.hadoop.fs.ChecksumFileSystem.open(ChecksumFileSystem.java:346)
    at org.apache.hadoop.fs.FileSystem.open(FileSystem.java:769)
    at org.apache.hadoop.mapred.LineRecordReader.<init>(LineRecordReader.java:109)
    at org.apache.hadoop.mapred.TextInputFormat.getRecordReader(TextInputFormat.java:67)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.liftedTree1$1(HadoopRDD.scala:267)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.<init>(HadoopRDD.scala:266)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:224)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:95)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

[ERROR] Task 0 in stage 0.0 failed 4 times; aborting job
Exception in thread "main" org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.3 in stage 0.0 (TID 3, 10.128.16.2, executor 0): java.io.FileNotFoundException: File file:/C:/Users/Utente/cache_pretrained/explain_document_ml_en_2.4.0_2.4_1580252705962/metadata/part-00000 does not exist
    at org.apache.hadoop.fs.RawLocalFileSystem.deprecatedGetFileStatus(RawLocalFileSystem.java:611)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileLinkStatusInternal(RawLocalFileSystem.java:824)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileStatus(RawLocalFileSystem.java:601)
    at org.apache.hadoop.fs.FilterFileSystem.getFileStatus(FilterFileSystem.java:421)
    at org.apache.hadoop.fs.ChecksumFileSystem$ChecksumFSInputChecker.<init>(ChecksumFileSystem.java:142)
    at org.apache.hadoop.fs.ChecksumFileSystem.open(ChecksumFileSystem.java:346)
    at org.apache.hadoop.fs.FileSystem.open(FileSystem.java:769)
    at org.apache.hadoop.mapred.LineRecordReader.<init>(LineRecordReader.java:109)
    at org.apache.hadoop.mapred.TextInputFormat.getRecordReader(TextInputFormat.java:67)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.liftedTree1$1(HadoopRDD.scala:267)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.<init>(HadoopRDD.scala:266)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:224)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:95)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

Driver stacktrace:
    at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1891)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1879)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1878)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
    at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1878)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:927)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:927)
    at scala.Option.foreach(Option.scala:257)
    at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:927)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2112)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2061)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2050)
    at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
    at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:738)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2061)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2082)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2101)
    at org.apache.spark.rdd.RDD$$anonfun$take$1.apply(RDD.scala:1409)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
    at org.apache.spark.rdd.RDD.withScope(RDD.scala:385)
    at org.apache.spark.rdd.RDD.take(RDD.scala:1382)
    at org.apache.spark.rdd.RDD$$anonfun$first$1.apply(RDD.scala:1423)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
    at org.apache.spark.rdd.RDD.withScope(RDD.scala:385)
    at org.apache.spark.rdd.RDD.first(RDD.scala:1422)
    at org.apache.spark.ml.util.DefaultParamsReader$.loadMetadata(ReadWrite.scala:615)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$.load(Pipeline.scala:267)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:348)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:342)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:374)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:368)
    at com.johnsnowlabs.nlp.pretrained.PretrainedPipeline.<init>(PretrainedPipeline.scala:26)
    at com.johnsnowlabs.nlp.pretrained.PretrainedPipeline.<init>(PretrainedPipeline.scala:21)
    at tags_extraction.tags_extraction_eng$.main(tags_extraction_eng.scala:157)
    at tags_extraction.tags_extraction_eng.main(tags_extraction_eng.scala)
Caused by: java.io.FileNotFoundException: File file:/C:/Users/Utente/cache_pretrained/explain_document_ml_en_2.4.0_2.4_1580252705962/metadata/part-00000 does not exist
    at org.apache.hadoop.fs.RawLocalFileSystem.deprecatedGetFileStatus(RawLocalFileSystem.java:611)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileLinkStatusInternal(RawLocalFileSystem.java:824)
    at org.apache.hadoop.fs.RawLocalFileSystem.getFileStatus(RawLocalFileSystem.java:601)
    at org.apache.hadoop.fs.FilterFileSystem.getFileStatus(FilterFileSystem.java:421)
    at org.apache.hadoop.fs.ChecksumFileSystem$ChecksumFSInputChecker.<init>(ChecksumFileSystem.java:142)
    at org.apache.hadoop.fs.ChecksumFileSystem.open(ChecksumFileSystem.java:346)
    at org.apache.hadoop.fs.FileSystem.open(FileSystem.java:769)
    at org.apache.hadoop.mapred.LineRecordReader.<init>(LineRecordReader.java:109)
    at org.apache.hadoop.mapred.TextInputFormat.getRecordReader(TextInputFormat.java:67)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.liftedTree1$1(HadoopRDD.scala:267)
    at org.apache.spark.rdd.HadoopRDD$$anon$1.<init>(HadoopRDD.scala:266)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:224)
    at org.apache.spark.rdd.HadoopRDD.compute(HadoopRDD.scala:95)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

This is not a pipeline issue rather the environment configuration and how it uses the HDFS.
When you use this command PretrainedPipeline (" recognize_entities_dl ", lang =" en ") it works, it downloads it on an available file system which in your case is local not HDFS because your Apache Spark is not well set to use HDFS when it's up and running. You can also use the downloaded/extracted pipeline by pointing to it on your local file system.

You cannot load a pre-trained pipeline from an HDFS when your environment is not set up correctly to have HDFS. Some pipelines will use the same path to extract some indexes on HDFS but in your setup, HDFS is not a default file system.

NOTE: If the cluster is configured correctly, then HDFS will be the only file system. There is no need to use IP:PORT and there is no way to access any local path as everything with or without hdfs:/// will be assumed HDFS.

I'll close this as this is not an issue with Spark NLP and more about Apache Spark + HDFS in a none Hadoop setup. (please feel free to continue the discussion if there is something new)

Hi, I managed to set hdfs as the default file system. Now when I try to download the pipeline the zip file is created on hdfs in the "user/root/cache-pretrained" folder. However the download is not successful and this error is thrown (the same happens with "explain_document_ml"):

Exception in thread "main" java.lang.IllegalArgumentException: requirement failed: Was not found appropriate resource to download for request: ResourceRequest(recognize_entities_dl,Some(en),public/models,2.4.5,2.4.5) with downloader: com.johnsnowlabs.nlp.pretrained.S3ResourceDownloader@7c2924d7
    at scala.Predef$.require(Predef.scala:224)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadResource(ResourceDownloader.scala:341)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:372)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:367)
    at com.johnsnowlabs.nlp.pretrained.PretrainedPipeline.<init>(PretrainedPipeline.scala:26)
    at com.johnsnowlabs.nlp.pretrained.PretrainedPipeline.<init>(PretrainedPipeline.scala:21)
    at tags_extraction.tags_extraction_eng$.main(tags_extraction_eng.scala:171)
    at tags_extraction.tags_extraction_eng.main(tags_extraction_eng.scala)

Here is the updated code:

import com.johnsnowlabs.nlp.pretrained.PretrainedPipeline
import com.johnsnowlabs.util.{ConfigLoader, PipelineModels}
import org.apache.hadoop.fs.FileSystem
import org.apache.spark.SparkContext
import org.apache.spark.ml.PipelineModel
import org.apache.spark.sql.{DataFrame, SparkSession}
import utils.Schemas.feedSchema

object tags_extraction_eng {

  System.setProperty("HADOOP_USER_NAME", "root")

  val (spark, sc) = start_spark_session()

  val fs = FileSystem.get(sc.hadoopConfiguration)
  println("HDFS_IMPL:", sc.hadoopConfiguration.get("fs.hdfs.impl"))  //null
  println("DEFAULT_FS:", sc.hadoopConfiguration.get("fs.defaultFS"))  // hdfs://remote_ip:remote_port
  println("FS_URI:", fs.getUri.toString)  // hdfs://remote_ip:remote_port
  println("FS_SCHEME:", fs.getScheme)  // hdfs


  def start_spark_session(): (SparkSession, SparkContext) = {
    val spark = SparkSession.builder()
      .appName("feeds_tags_extraction_nlp_eng")
      //            .master("local[*]")
      .master("spark://remote_ip:7077")
      .config("spark.driver.host", "host_ip")

      .config("spark.submit.deployMode", "client")

 .config("sparknlp.settings.pretrained.cache_folder","hdfs://remote_ip:remote_port/user/root/cache_pretrained")
      .config("sparknlp.settings.storage.hadoop.tmp.dir", "hdfs://remote_ip:remote_port/tmp/hadoop-root")

      .config("spark.jars", 
        "https://repo1.maven.org/maven2/com/johnsnowlabs/nlp/spark-nlp_2.11/2.4.5/spark-nlp_2.11-2.4.5.jar," +       
        "hdfs://remote_ip:remote_port/spark_project.jar"
      )

    spark.sparkContext.hadoopConfiguration.set("fs.defaultFS", "hdfs://remote_ip:remote_port")
    spark.sparkContext.hadoopConfiguration.set("hadoop.http.staticuser.user", "root")
    spark.sparkContext.hadoopConfiguration.set("hadoop.proxyuser.hue.hosts", "*")
    spark.sparkContext.hadoopConfiguration.set("hadoop.proxyuser.hue.groups", "*")
    spark.sparkContext.hadoopConfiguration.set("io.compression.codecs", "org.apache.hadoop.io.compress.SnappyCodec")
    spark.sparkContext.hadoopConfiguration.set("dfs.client.use.datanode.hostname", "true")
    spark.sparkContext.hadoopConfiguration.set("dfs.datanode.use.datanode.hostname", "true")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.datanode.registration.ip-hostname-check", "false")
    spark.sparkContext.hadoopConfiguration.set("dfs.permissions.enabled", "false")
    spark.sparkContext.hadoopConfiguration.set("dfs.replication", "1")
    spark.sparkContext.hadoopConfiguration.set("hadoop.tmp.dir", "/tmp/hadoop-root")

    val sc = spark.sparkContext
    sc.getConf.getAll.foreach(println)

    (spark, sc)
  }


  def main(args: Array[String]) {
    val pipeline = new PretrainedPipeline(downloadName = "recognize_entities_dl", lang = "en")
    print(pipeline)

    spark.close()
  }

}

I am glad the HDSF part is resolved.

A few notes:

  • When a model or a pipeline is not found it's either the name has an issue or there is a connection problem due to internet connectivity or firewall/proxy.
  • The JAR you are using directly from the maven is not a Fat jar. Please either use packages com.johnsnowlabs.nlp:spark-nlp_2.11:2.5.0 or download the Fat JAR from the release notes or use the direct link to S3 and use the same spark.jars
  • The correct command is val pipeline = PretrainedPipeline("recognize_entities_dl", lang="en") there is no new or downloadName. It's actually name.

Could you please check the internet, use the S3 link instead of Maven for spark.jars and correct the command and try again?

I have added to spark.jars:

and modified the command as you suggests, but I keep getting:

[ WARN] DataStreamer Exception
java.nio.channels.UnresolvedAddressException
    at sun.nio.ch.Net.checkAddress(Net.java:101)
    at sun.nio.ch.SocketChannelImpl.connect(SocketChannelImpl.java:622)
    at org.apache.hadoop.net.SocketIOWithTimeout.connect(SocketIOWithTimeout.java:192)
    at org.apache.hadoop.net.NetUtils.connect(NetUtils.java:530)
    at org.apache.hadoop.hdfs.DFSOutputStream.createSocketForPipeline(DFSOutputStream.java:1606)
    at org.apache.hadoop.hdfs.DFSOutputStream$DataStreamer.createBlockOutputStream(DFSOutputStream.java:1404)
    at org.apache.hadoop.hdfs.DFSOutputStream$DataStreamer.nextBlockOutputStream(DFSOutputStream.java:1357)
    at org.apache.hadoop.hdfs.DFSOutputStream$DataStreamer.run(DFSOutputStream.java:587)
Exception in thread "main" java.lang.IllegalArgumentException: requirement failed: Was not found appropriate resource to download for request: ResourceRequest(explain_document_ml,Some(en),public/models,2.5.0,2.4.5) with downloader: com.johnsnowlabs.nlp.pretrained.S3ResourceDownloader@297c9a9b
    at scala.Predef$.require(Predef.scala:224)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadResource(ResourceDownloader.scala:342)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:373)
    at com.johnsnowlabs.nlp.pretrained.ResourceDownloader$.downloadPipeline(ResourceDownloader.scala:368)
    at com.johnsnowlabs.nlp.pretrained.PretrainedPipeline.<init>(PretrainedPipeline.scala:26)
    at tags_extraction.tags_extraction_eng$.main(tags_extraction_eng.scala:174)
    at tags_extraction.tags_extraction_eng.main(tags_extraction_eng.scala)

Below I attach a screenshot of the screen on HDFS:

Annotazione 2020-05-24 134545

Why do I keep seeing explain_document_ml when your code is trying to use another pipeline?

In any case, if one or all the pre-trained models/pipelines cannot be downloaded, you can always download them manually, extract them, and try to load them instead of pretained():

https://github.com/JohnSnowLabs/spark-nlp-models

PS: Do you have any AWS config on your machine? Like a profile or any sort of access to S3? (There could be a conflict between the two if there is one)

PS2: I can see it downloaded the correct pipeline, there must be a permission issue that it couldn't extract it and it fails with "cannot find the pipeline". Could you also check the permission of that directory to be sure your use (the one creates the SparkSession) has full permission over that path?

  • Why do I keep seeing explain_document_ml when your code is trying to use another pipeline?
    Because I was trying with another pipeline

  • In any case, if one or all the pre-trained models/pipelines cannot be downloaded, you can always download them manually, extract them, and try to load them instead of pretained()
    I have already tried to download the pipeline manually from HDFS, but
    as soon as it creates the folder_cdx in "tmp/hadoop-root" folder it fails, show the following error and delete the folder_cdx (loading fails when reading storge folder in stage 3 of the pipeline):

Exception in thread "main" java.lang.ExceptionInInitializerError
    at hdfs.hdfs_connect.main(hdfs_connect.scala)
Caused by: java.nio.channels.UnresolvedAddressException
    at sun.nio.ch.Net.checkAddress(Net.java:101)
    at sun.nio.ch.SocketChannelImpl.connect(SocketChannelImpl.java:622)
    at org.apache.hadoop.net.SocketIOWithTimeout.connect(SocketIOWithTimeout.java:192)
    at org.apache.hadoop.net.NetUtils.connect(NetUtils.java:530)
    at org.apache.hadoop.hdfs.DFSClient.newConnectedPeer(DFSClient.java:3090)
    at org.apache.hadoop.hdfs.BlockReaderFactory.nextTcpPeer(BlockReaderFactory.java:778)
    at org.apache.hadoop.hdfs.BlockReaderFactory.getRemoteBlockReaderFromTcp(BlockReaderFactory.java:693)
    at org.apache.hadoop.hdfs.BlockReaderFactory.build(BlockReaderFactory.java:354)
    at org.apache.hadoop.hdfs.DFSInputStream.blockSeekTo(DFSInputStream.java:617)
    at org.apache.hadoop.hdfs.DFSInputStream.readWithStrategy(DFSInputStream.java:841)
    at org.apache.hadoop.hdfs.DFSInputStream.read(DFSInputStream.java:889)
    at java.io.DataInputStream.read(DataInputStream.java:100)
    at org.apache.hadoop.io.IOUtils.copyBytes(IOUtils.java:78)
    at org.apache.hadoop.io.IOUtils.copyBytes(IOUtils.java:52)
    at org.apache.hadoop.io.IOUtils.copyBytes(IOUtils.java:112)
    at org.apache.hadoop.fs.FileUtil.copy(FileUtil.java:366)
    at org.apache.hadoop.fs.FileUtil.copy(FileUtil.java:356)
    at org.apache.hadoop.fs.FileUtil.copy(FileUtil.java:338)
    at com.johnsnowlabs.storage.StorageHelper$.copyIndexToCluster(StorageHelper.scala:66)
    at com.johnsnowlabs.storage.StorageHelper$.sendToCluster(StorageHelper.scala:53)
    at com.johnsnowlabs.storage.StorageHelper$.load(StorageHelper.scala:27)
    at com.johnsnowlabs.storage.HasStorageModel$$anonfun$deserializeStorage$1.apply(HasStorageModel.scala:27)
    at com.johnsnowlabs.storage.HasStorageModel$$anonfun$deserializeStorage$1.apply(HasStorageModel.scala:26)
    at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
    at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
    at com.johnsnowlabs.storage.HasStorageModel$class.deserializeStorage(HasStorageModel.scala:26)
    at com.johnsnowlabs.nlp.embeddings.WordEmbeddingsModel.deserializeStorage(WordEmbeddingsModel.scala:32)
    at com.johnsnowlabs.storage.StorageReadable$class.readStorage(StorageReadable.scala:23)
    at com.johnsnowlabs.nlp.embeddings.WordEmbeddingsModel$.readStorage(WordEmbeddingsModel.scala:156)
    at com.johnsnowlabs.storage.StorageReadable$$anonfun$1.apply(StorageReadable.scala:26)
    at com.johnsnowlabs.storage.StorageReadable$$anonfun$1.apply(StorageReadable.scala:26)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead$1.apply(ParamsAndFeaturesReadable.scala:31)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead$1.apply(ParamsAndFeaturesReadable.scala:30)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$class.com$johnsnowlabs$nlp$ParamsAndFeaturesReadable$$onRead(ParamsAndFeaturesReadable.scala:30)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$read$1.apply(ParamsAndFeaturesReadable.scala:41)
    at com.johnsnowlabs.nlp.ParamsAndFeaturesReadable$$anonfun$read$1.apply(ParamsAndFeaturesReadable.scala:41)
    at com.johnsnowlabs.nlp.FeaturesReader.load(ParamsAndFeaturesReadable.scala:19)
    at com.johnsnowlabs.nlp.FeaturesReader.load(ParamsAndFeaturesReadable.scala:8)
    at org.apache.spark.ml.util.DefaultParamsReader$.loadParamsInstance(ReadWrite.scala:652)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$$anonfun$4.apply(Pipeline.scala:274)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$$anonfun$4.apply(Pipeline.scala:272)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
    at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
    at scala.collection.TraversableLike$class.map(TraversableLike.scala:234)
    at scala.collection.mutable.ArrayOps$ofRef.map(ArrayOps.scala:186)
    at org.apache.spark.ml.Pipeline$SharedReadWrite$.load(Pipeline.scala:272)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:348)
    at org.apache.spark.ml.PipelineModel$PipelineModelReader.load(Pipeline.scala:342)
    at org.apache.spark.ml.util.MLReadable$class.load(ReadWrite.scala:380)
    at org.apache.spark.ml.PipelineModel$.load(Pipeline.scala:332)
    at hdfs.hdfs_connect$.<init>(hdfs_connect.scala:13)
    at hdfs.hdfs_connect$.<clinit>(hdfs_connect.scala)
    ... 1 more

PS: Do you have any AWS config on your machine? Like a profile or any sort of access to S3? (There could be a conflict between the two if there is one)
No, I haven't

PS2: I can see it downloaded the correct pipeline, there must be a permission issue that it couldn't extract it and it fails with "cannot find the pipeline". Could you also check the permission of that directory to be sure your use (the one creates the SparkSession) has full permission over that path?
The screenshot shows the permissions associated with the file, however if I try to save files or pipelines created by me there are no problems, so I think it is not a problem of user permissions. Moreover I don't think it gets to the extraction phase because the file is empty, I also tried to download it manually from HDFS on my PC but when I try to open it it says that the file is corrupt or damaged.

I managed to load the downloaded pipeline manually, there was a host problem. However now I get the following error:

[Stage 22:>                                                         (0 + 1) / 1][ WARN] Lost task 0.0 in stage 22.0 (TID 42, 10.128.16.2, executor 0): java.lang.UnsatisfiedLinkError: /tmp/tensorflow_native_libraries-1590409637658-0/libtensorflow_jni.so: Error relocating /tmp/tensorflow_native_libraries-1590409637658-0/libtensorflow_jni.so: __memcpy_chk: symbol not found
    at java.lang.ClassLoader$NativeLibrary.load(Native Method)
    at java.lang.ClassLoader.loadLibrary0(ClassLoader.java:1946)
    at java.lang.ClassLoader.loadLibrary(ClassLoader.java:1828)
    at java.lang.Runtime.load0(Runtime.java:810)
    at java.lang.System.load(System.java:1088)
    at org.tensorflow.NativeLibrary.load(NativeLibrary.java:101)
    at org.tensorflow.TensorFlow.init(TensorFlow.java:67)
    at org.tensorflow.TensorFlow.<clinit>(TensorFlow.java:82)
    at org.tensorflow.Graph.<clinit>(Graph.java:479)
    at com.johnsnowlabs.ml.tensorflow.TensorflowWrapper.getSession(TensorflowWrapper.scala:56)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:66)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:56)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer.predict(TensorflowNer.scala:56)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.tag(NerDLModel.scala:135)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.annotate(NerDLModel.scala:188)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:35)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:34)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:255)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

[ WARN] Lost task 0.1 in stage 22.0 (TID 43, 10.128.16.2, executor 0): java.lang.NoClassDefFoundError: Could not initialize class org.tensorflow.Graph
    at com.johnsnowlabs.ml.tensorflow.TensorflowWrapper.getSession(TensorflowWrapper.scala:56)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:66)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:56)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer.predict(TensorflowNer.scala:56)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.tag(NerDLModel.scala:135)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.annotate(NerDLModel.scala:188)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:35)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:34)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:255)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

[ERROR] Task 0 in stage 22.0 failed 4 times; aborting job
Exception in thread "main" org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 22.0 failed 4 times, most recent failure: Lost task 0.3 in stage 22.0 (TID 45, 10.128.16.2, executor 0): java.lang.NoClassDefFoundError: Could not initialize class org.tensorflow.Graph
    at com.johnsnowlabs.ml.tensorflow.TensorflowWrapper.getSession(TensorflowWrapper.scala:56)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:66)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:56)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer.predict(TensorflowNer.scala:56)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.tag(NerDLModel.scala:135)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.annotate(NerDLModel.scala:188)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:35)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:34)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:255)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

Driver stacktrace:
    at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1891)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1879)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1878)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
    at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1878)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:927)
    at org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:927)
    at scala.Option.foreach(Option.scala:257)
    at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:927)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2112)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2061)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2050)
    at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
    at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:738)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2061)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2082)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:2101)
    at org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:365)
    at org.apache.spark.sql.execution.CollectLimitExec.executeCollect(limit.scala:38)
    at org.apache.spark.sql.Dataset.org$apache$spark$sql$Dataset$$collectFromPlan(Dataset.scala:3389)
    at org.apache.spark.sql.Dataset$$anonfun$head$1.apply(Dataset.scala:2550)
    at org.apache.spark.sql.Dataset$$anonfun$head$1.apply(Dataset.scala:2550)
    at org.apache.spark.sql.Dataset$$anonfun$52.apply(Dataset.scala:3370)
    at org.apache.spark.sql.execution.SQLExecution$$anonfun$withNewExecutionId$1.apply(SQLExecution.scala:80)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:127)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:75)
    at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3369)
    at org.apache.spark.sql.Dataset.head(Dataset.scala:2550)
    at org.apache.spark.sql.Dataset.take(Dataset.scala:2764)
    at org.apache.spark.sql.Dataset.getRows(Dataset.scala:254)
    at org.apache.spark.sql.Dataset.showString(Dataset.scala:291)
    at org.apache.spark.sql.Dataset.show(Dataset.scala:751)
    at org.apache.spark.sql.Dataset.show(Dataset.scala:710)
    at org.apache.spark.sql.Dataset.show(Dataset.scala:719)
    at tags_extraction.tags_extraction_eng$.main(tags_extraction_eng.scala:189)
    at tags_extraction.tags_extraction_eng.main(tags_extraction_eng.scala)
Caused by: java.lang.NoClassDefFoundError: Could not initialize class org.tensorflow.Graph
    at com.johnsnowlabs.ml.tensorflow.TensorflowWrapper.getSession(TensorflowWrapper.scala:56)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:66)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer$$anonfun$predict$2.apply(TensorflowNer.scala:56)
    at scala.collection.Iterator$class.foreach(Iterator.scala:891)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)
    at com.johnsnowlabs.ml.tensorflow.TensorflowNer.predict(TensorflowNer.scala:56)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.tag(NerDLModel.scala:135)
    at com.johnsnowlabs.nlp.annotators.ner.dl.NerDLModel.annotate(NerDLModel.scala:188)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:35)
    at com.johnsnowlabs.nlp.AnnotatorModel$$anonfun$dfAnnotate$1.apply(AnnotatorModel.scala:34)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$13$$anon$1.hasNext(WholeStageCodegenExec.scala:636)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:255)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:858)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:310)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

This is the updated code:

import com.johnsnowlabs.nlp.pretrained.PretrainedPipeline
import com.johnsnowlabs.util.{ConfigLoader, PipelineModels}
import org.apache.hadoop.fs.FileSystem
import org.apache.spark.SparkContext
import org.apache.spark.ml.PipelineModel
import org.apache.spark.sql.{DataFrame, SparkSession}
import utils.Schemas.feedSchema

object tags_extraction_eng {

  System.setProperty("HADOOP_USER_NAME", "root")

  val (spark, sc) = start_spark_session()

  val fs = FileSystem.get(sc.hadoopConfiguration)
  println("HDFS_IMPL:", sc.hadoopConfiguration.get("fs.hdfs.impl"))  //null
  println("DEFAULT_FS:", sc.hadoopConfiguration.get("fs.defaultFS"))  // hdfs://remote_ip:remote_port
  println("FS_URI:", fs.getUri.toString)  // hdfs://remote_ip:remote_port
  println("FS_SCHEME:", fs.getScheme)  // hdfs


  def start_spark_session(): (SparkSession, SparkContext) = {
    val spark = SparkSession.builder()
      .appName("feeds_tags_extraction_nlp_eng")
      //            .master("local[*]")
      .master("spark://remote_ip:7077")
      .config("spark.driver.host", "host_ip")

      .config("spark.submit.deployMode", "client")

 .config("sparknlp.settings.pretrained.cache_folder","hdfs://remote_ip:remote_port/user/root/cache_pretrained")
      .config("sparknlp.settings.storage.hadoop.tmp.dir", "hdfs://remote_ip:remote_port/tmp/hadoop-root")

      .config("spark.jars", "hdfs://192.168.150.71:9000/sparkscala.jar," +
        "https://repo1.maven.org/maven2/com/couchbase/client/java-client/2.7.6/java-client-2.7.6.jar," +
        "https://repo1.maven.org/maven2/com/couchbase/client/core-io/1.7.6/core-io-1.7.6.jar," +
        "https://repo1.maven.org/maven2/io/reactivex/rxjava/1.3.8/rxjava-1.3.8.jar," +
        "https://repo1.maven.org/maven2/io/reactivex/rxscala_2.12/0.26.5/rxscala_2.12-0.26.5.jar," +
        "https://repo1.maven.org/maven2/io/opentracing/opentracing-api/0.31.0/opentracing-api-0.31.0.jar," +
        "https://repo1.maven.org/maven2/com/couchbase/client/spark-connector_2.11/2.3.0/spark-connector_2.11-2.3.0.jar," +
        "https://s3.amazonaws.com/auxdata.johnsnowlabs.com/public/spark-nlp-assembly-2.5.0.jar," +
        "https://repo1.maven.org/maven2/org/tensorflow/tensorflow/1.15.0/tensorflow-1.15.0.jar," +
        "https://repo1.maven.org/maven2/org/tensorflow/libtensorflow/1.15.0/libtensorflow-1.15.0.jar," +
        "https://repo1.maven.org/maven2/org/tensorflow/libtensorflow_jni/1.15.0/libtensorflow_jni-1.15.0.jar," +
        "https://repo1.maven.org/maven2/org/tensorflow/tensorflow-hadoop/1.15.0/tensorflow-hadoop-1.15.0.jar," +
        "https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-common/3.2.1/hadoop-common-3.2.1.jar,"
      )

    spark.sparkContext.hadoopConfiguration.set("fs.defaultFS", "hdfs://remote_ip:remote_port")
    spark.sparkContext.hadoopConfiguration.set("hadoop.http.staticuser.user", "root")
    spark.sparkContext.hadoopConfiguration.set("hadoop.proxyuser.hue.hosts", "*")
    spark.sparkContext.hadoopConfiguration.set("hadoop.proxyuser.hue.groups", "*")
    spark.sparkContext.hadoopConfiguration.set("io.compression.codecs", "org.apache.hadoop.io.compress.SnappyCodec")
    spark.sparkContext.hadoopConfiguration.set("dfs.client.use.datanode.hostname", "true")
    spark.sparkContext.hadoopConfiguration.set("dfs.datanode.use.datanode.hostname", "true")
    spark.sparkContext.hadoopConfiguration.set("dfs.webhdfs.enabled", "true")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.datanode.registration.ip-hostname-check", "false")
    spark.sparkContext.hadoopConfiguration.set("dfs.permissions.enabled", "false")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.rpc-bind-host", "0.0.0.0")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.servicerpc-bind-host", "0.0.0.0")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.http-bind-host", "0.0.0.0")
    spark.sparkContext.hadoopConfiguration.set("dfs.namenode.https-bind-host", "0.0.0.0")
    spark.sparkContext.hadoopConfiguration.set("dfs.replication", "1")
    spark.sparkContext.hadoopConfiguration.set("hadoop.tmp.dir", "/tmp/hadoop-root/")

    val sc = spark.sparkContext
    sc.getConf.getAll.foreach(println)

    (spark, sc)
  }


  def main(args: Array[String]) {
    val pipeline = PipelineModel.load("recognize_entities_dl_en_2.4.3_2.4_1584626752821")
    pipeline.transform(df)

    spark.close()
  }

}

Hi this is a new issue, could you please create a new issue? Please fill the template correctly in details so we can reproduce. (Specially the OS, environment etc.)

Your SparkSession has changed, why are you adding all those dependencies manually and what is sparkscala? Where is spark-nlp package or fat jar?
Please follow the instructions as how to install/use spark-nlp vis either packages or fat jar. Also, you don鈥檛 need to add all those manually, especially the tensorflow.

what is sparkscala?
sparkscal.jar contains my codes;

Where is spark-nlp package or fat jar?
spark fat jar is in "https://s3.amazonaws.com/auxdata.johnsnowlabs.com/public/spark-nlp-ssembly-2.5.0.jar" (I found only this link to fat jar);

You don鈥檛 need to add all those manually, especially the tensorflow.
Yes, I know it. I only did one test since I had seen the NoClassDefFoundError error.

Anyway I just created a new issue (with environment info, code updated and errors).

Was this page helpful?
0 / 5 - 0 ratings