【问题标题】:Spark in cluster with Docker: BlockManagerId; local class incompatibleSpark 在集群中使用 Docker:BlockManagerId;本地类不兼容
【发布时间】:2018-12-01 09:22:20
【问题描述】:

我在使用 Spark 和 Docker 分发操作时遇到类型不匹配的问题。 The tutorial我关注的似乎很清楚。这是我对 Scala 代码的尝试:

package test

import com.datastax.spark.connector.cql.CassandraConnector
import org.apache.spark.{SparkConf, SparkContext}
import readhub.sharedkernel.config.Settings

object Application extends App {
    import com.datastax.spark.connector._


    val conf = new SparkConf(true)
      .setAppName("Coordinator")
      .setMaster("spark://localhost:7077")
      .set("spark.cassandra.connection.host", "valid host")

    val sc = new SparkContext(conf)

    CassandraConnector(conf).withSessionDo { session =>
      session.execute("CREATE KEYSPACE test2 WITH REPLICATION = {'class': 'SimpleStrategy', 'replication_factor': 1 }")
      session.execute("CREATE TABLE test2.words (word text PRIMARY KEY, count int)")
      session.execute("INSERT INTO test2.words(word, count) VALUES('hey', 32)")

      sc.cassandraTable("test2", "words")
        .map(r => r.getString("word"))
        .foreach(process)

    }

  def process(word: String): Unit = {
    // Dummy processing
    println(word)
  }
} 

build.sbt 看起来像这样:

import sbt.project

val sparkSql = "org.apache.spark" %% "spark-sql" % "2.3.0" % "provided"
val sparkCassandraConnector = "com.datastax.spark" %% "spark-cassandra-connector" % "2.3.0" % "provided"

lazy val commonSettings = Seq(
  version := "0.1",
  scalaVersion := "2.11.12",
  organization := "ch.heig-vd"
)

lazy val root = (project in file("."))
  .settings(
    commonSettings,
    name := "Root"
  )
  .aggregate(
    coordinator
  )

lazy val coordinator = project
  .settings(
    commonSettings,
    name := "Coordinator",
    libraryDependencies ++= Seq(
      sparkSql,
      sparkCassandraConnector
    )
  )

Dockerfile 取自 this image 并稍作修改以使用 Spark 2.3.0 版本:

FROM phusion/baseimage:0.9.22

ENV SPARK_VERSION 2.3.0
ENV SPARK_INSTALL /usr/local
ENV SPARK_HOME $SPARK_INSTALL/spark
ENV SPARK_ROLE master
ENV HADOOP_VERSION 2.7
ENV SPARK_MASTER_PORT 7077
ENV PYSPARK_PYTHON python3
ENV DOCKERIZE_VERSION v0.2.0

RUN apt-get update && \
    apt-get install -y openjdk-8-jdk autossh python3-pip && \
    apt-get clean && \
    rm -rf /var/lib/apt/lists/* /tmp/* /var/tmp/*

##### INSTALL DOCKERIZE
RUN curl -L -O https://github.com/jwilder/dockerize/releases/download/$DOCKERIZE_VERSION/dockerize-linux-amd64-$DOCKERIZE_VERSION.tar.gz && \
    tar -C /usr/local/bin -xzvf dockerize-linux-amd64-$DOCKERIZE_VERSION.tar.gz && \
    rm -rf dockerize-linux-amd64-$DOCKERIZE_VERSION.tar.gz

##### INSTALL APACHE SPARK WITH HDFS
RUN curl -s http://mirror.synyx.de/apache/spark/spark-$SPARK_VERSION/spark-$SPARK_VERSION-bin-hadoop$HADOOP_VERSION.tgz | tar -xz -C $SPARK_INSTALL && \
    cd $SPARK_INSTALL && ln -s spark-$SPARK_VERSION-bin-hadoop$HADOOP_VERSION spark

WORKDIR $SPARK_HOME

##### ADD Scripts
RUN mkdir /etc/service/spark
ADD runit/spark.sh /etc/service/spark/run
RUN chmod +x /etc/service/**/*

EXPOSE 4040 6066 7077 7078 8080 8081 8888

VOLUME ["$SPARK_HOME/logs"]

CMD ["/sbin/my_init"]

docker-compose.yml 也很简单:

version: "3"

services:
  master:
    build: birgerk-apache-spark

    ports:
      - "7077:7077"
      - "8080:8080"

  slave:
    build: birgerk-apache-spark
    environment:
      - SPARK_ROLE=slave
      - SPARK_MASTER=master
    depends_on:
      - master

我将 git repo 克隆到文件夹 birgerk-apache-spark 中,仅将 Spark 的版本更改为 2.3.0。

最后,我使用以下方法粘合所有内容:

sbt coordinator/assembly

创建胖罐子和

spark-submit --class test.Application --packages com.datastax.spark:spark-cassandra-connector_2.11:2.3.0 --master spark://localhost:7077 ReadHub\ Coordinator-assembly-0.1.jar

将 jar 提交到集群中。当我发出spark-submit 时出现错误:

ERROR TransportRequestHandler:199 - 调用时出错 RPC id 7068633004064450609 上的 RpcHandler#receive() java.io.InvalidClassException: org.apache.spark.storage.BlockManagerId;本地类不兼容: stream classdesc serialVersionUID = 6155820641931972169,本地类 serialVersionUID = -3720498261147521051 在 java.io.ObjectStreamClass.initNonProxy(ObjectStreamClass.java:687) 在 java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1876) 在 java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1745) 在 java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2033) 在 java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1567) 在 java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2278) 在 java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2202) [..]

从我的角度来看,Dockerfile 正确下载了相应版本的 Spark,可以在我的 build.sbt 中作为依赖项找到该版本。

我确定我错过了一些基本的东西。谁能指出我正确的方向?

非常感谢!

【问题讨论】:

  • 我只在运行不匹配版本的 spark 时看到该错误。我会仔细检查您的 docker 映像版本是否是实际运行的版本。
  • 我对@9​​87654331@ 不太熟悉,但是在将目录传递给build: 指令时,文档显示了一个领先的./。在您的撰写文件中尝试build: ./birgerk-apache-spark
  • 您好@TravisHegner,感谢您的评论!我通过浏览 localhost:8080 检查了 Spark 的版本,该版本是预期的版本(2.3.0)。我还更改了build 指令并发出了docker-compose build,但没有更改。 Docker 正确地定位了构建文件夹。我也有一种不匹配的感觉,但我找不到它。 :(
  • 也许您的应用程序中的其他一些依赖项正在引入不匹配的 spark 版本?你的 Java 版本呢,它与容器中的版本匹配吗?

标签: scala apache-spark docker cluster-computing


【解决方案1】:

spark 2.3.3 和 spark 2.3.0 之间的版本不匹配。

注意不要提交在您的主机上定义了 SPARK_HOME 的作业,这可能会导致此类问题

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-06-01
    • 2016-11-28
    • 2018-07-20
    • 2020-01-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多