java.lang.NoSuchMethodError: org.apache.kafka.clients.consumer.KafkaConsumer.assign

1/2/2019

My spark job running well when submitted to local spark cluster (spark-2.3.1-bin-hadoop2.7). But I got several errors when submitting it to spark cluster (spark-2.3.1-bin-hadoop3.0.0) on k8s cluster on openstack.

At beginning, I got Exception in thread "main" org.apache.kafka.common.config.ConfigException: Missing required configuration "partition.assignment.strategy" which has no default value.. so, I added the following configuration:

"partition.assignment.strategy" -> "org.apache.kafka.clients.consumer.RangeAssignor", 

Now, I got the error:

Exception in thread "streaming-start" java.lang.NoSuchMethodError: org.apache.kafka.clients.consumer.KafkaConsumer.assign(Ljava/util/Collection;)V
    at org.apache.spark.streaming.kafka010.Assign.onStart(ConsumerStrategy.scala:193)
    at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.consumer(DirectKafkaInputDStream.scala:70)
    at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.start(DirectKafkaInputDStream.scala:240)
    at org.apache.spark.streaming.DStreamGraph$anonfun$start$7.apply(DStreamGraph.scala:54)
    at org.apache.spark.streaming.DStreamGraph$anonfun$start$7.apply(DStreamGraph.scala:54)
    at scala.collection.parallel.mutable.ParArray$ParArrayIterator.foreach_quick(ParArray.scala:143)
    at scala.collection.parallel.mutable.ParArray$ParArrayIterator.foreach(ParArray.scala:136)
    at scala.collection.parallel.ParIterableLike$Foreach.leaf(ParIterableLike.scala:972)
    at scala.collection.parallel.Task$anonfun$tryLeaf$1.apply$mcV$sp(Tasks.scala:49)
    at scala.collection.parallel.Task$anonfun$tryLeaf$1.apply(Tasks.scala:48)
    at scala.collection.parallel.Task$anonfun$tryLeaf$1.apply(Tasks.scala:48)
    at scala.collection.parallel.Task$class.tryLeaf(Tasks.scala:51)
    at scala.collection.parallel.ParIterableLike$Foreach.tryLeaf(ParIterableLike.scala:969)
    at scala.collection.parallel.AdaptiveWorkStealingTasks$WrappedTask$class.compute(Tasks.scala:152)
    at scala.collection.parallel.AdaptiveWorkStealingForkJoinTasks$WrappedTask.compute(Tasks.scala:443)
    at scala.concurrent.forkjoin.RecursiveAction.exec(RecursiveAction.java:160)
    at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
    at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
    at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
    at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)
19/01/02 22:32:43 INFO ContextHandler: Started o.s.j.s.ServletContextHandler@3e8799f{/streaming,null,AVAILABLE,@Spark}
19/01/02 22:32:43 INFO ContextHandler: Started o.s.j.s.ServletContextHandler@3b353651{/streaming/json,null,AVAILABLE,@Spark}
19/01/02 22:32:43 INFO ContextHandler: Started o.s.j.s.ServletContextHandler@53d9826f{/streaming/batch,null,AVAILABLE,@Spark}
19/01/02 22:32:43 INFO ContextHandler: Started o.s.j.s.ServletContextHandler@1e84f3c8{/streaming/batch/json,null,AVAILABLE,@Spark}
19/01/02 22:32:43 INFO ContextHandler: Started o.s.j.s.ServletContextHandler@64f9f455{/static/streaming,null,AVAILABLE,@Spark}
19/01/02 22:32:43 INFO StreamingContext: StreamingContext started
19/01/02 22:32:43 INFO CoarseGrainedSchedulerBackend$DriverEndpoint: Registered executor NettyRpcEndpointRef(spark-client://Executor) (10.233.88.4:51780) with ID 0
19/01/02 22:32:43 INFO CoarseGrainedSchedulerBackend$DriverEndpoint: Registered executor NettyRpcEndpointRef(spark-client://Executor) (10.233.86.5:55632) with ID 1
19/01/02 22:32:43 INFO BlockManagerMasterEndpoint: Registering block manager 10.233.88.4:37243 with 366.3 MB RAM, BlockManagerId(0, 10.233.88.4, 37243, None)
19/01/02 22:32:43 INFO BlockManagerMasterEndpoint: Registering block manager 10.233.86.5:36954 with 366.3 MB RAM, BlockManagerId(1, 10.233.86.5, 36954, None)
19/01/02 22:32:44 INFO CoarseGrainedSchedulerBackend$DriverEndpoint: Registered executor NettyRpcEndpointRef(spark-client://Executor) (10.233.96.4:36610) with ID 2
19/01/02 22:32:44 INFO BlockManagerMasterEndpoint: Registering block manager 10.233.96.4:33930 with 366.3 MB RAM, BlockManagerId(2, 10.233.96.4, 33930, None)

and my spark job hangs here.

my pom.xml:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
  <modelVersion>4.0.0</modelVersion>
  <groupId>com.poc</groupId>
  <artifactId>spark-streaming</artifactId>
  <packaging>jar</packaging>
  <version>1.0-SNAPSHOT</version>
  <name>spark-streaming</name>
  <url>http://maven.apache.org</url>
  <properties>
    <java.version>1.8</java.version>
    <scala.version>2.11.8</scala.version>
    <scala.compactVersion>2.11</scala.compactVersion>
    <spark.version>2.3.1</spark.version>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
  </properties>
  <dependencies>
    <dependency>
      <groupId>junit</groupId>
      <artifactId>junit</artifactId>
      <version>3.8.1</version>
      <scope>test</scope>
    </dependency>
    <dependency>
      <groupId>org.scala-lang</groupId>
      <artifactId>scala-library</artifactId>
      <version>${scala.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-core_${scala.compactVersion}</artifactId>
      <version>${spark.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-sql_${scala.compactVersion}</artifactId>
      <version>${spark.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-streaming_${scala.compactVersion}</artifactId>
      <version>${spark.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.kafka</groupId>
      <artifactId>kafka-clients</artifactId>
      <version>0.10.0.0-SASL</version>
    </dependency>
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-streaming-kafka-0-10_${scala.compactVersion}</artifactId>
      <version>${spark.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-sql-kafka-0-10_${scala.compactVersion}</artifactId>
      <version>${spark.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.kafka</groupId>
      <artifactId>kafka_${scala.compactVersion}</artifactId>
      <version>0.10.0.0</version>
    </dependency>
    <dependency>
      <groupId>org.rogach</groupId>
      <artifactId>scallop_${scala.compactVersion}</artifactId>
      <version>2.1.2</version>
    </dependency>
  </dependencies>
  <pluginRepositories>
  </pluginRepositories>
  <build>
    <sourceDirectory>src/main/scala</sourceDirectory>
    <plugins>
      <plugin>
        <artifactId>maven-assembly-plugin</artifactId>
        <executions>
          <execution>
            <phase>package</phase>
            <goals>
              <goal>single</goal>
            </goals>
          </execution>
        </executions>
        <configuration>
          <archive>
            <manifest>
              <mainClass>com.poc.Job</mainClass>
            </manifest>
          </archive>
          <descriptorRefs>
            <descriptorRef>jar-with-dependencies</descriptorRef>
          </descriptorRefs>
        </configuration>
      </plugin>
      <plugin>
        <groupId>net.alchim31.maven</groupId>
        <artifactId>scala-maven-plugin</artifactId>
        <version>3.2.0</version>
        <executions>
          <execution>
            <id>scala-compile</id>
            <goals>
              <goal>compile</goal>
            </goals>
          </execution>
        </executions>
      </plugin>
    </plugins>
  </build>
</project>

Any hints welcomed. Thanks

-- BAE
apache-kafka
apache-spark
kubernetes
scala
spark-streaming-kafka

0 Answers