【问题标题】:Why does spark-submit fail to find kafka data source unless --packages is used?为什么 spark-submit 找不到 kafka 数据源,除非使用 --packages?
【发布时间】:2018-02-10 14:21:21
【问题描述】:

我正在尝试将 Kafka 集成到我的 Spark 应用程序中,这是我的 POM 文件所需的条目:

<dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
        <version>${spark.stream.kafka.version}</version>
</dependency>
<dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka_2.11</artifactId>
        <version>${kafka.version}</version>
</dependency>

对应的神器版本有:

<kafka.version>0.10.2.0</kafka.version>
<spark.stream.kafka.version>2.2.0</spark.stream.kafka.version>

我一直在摸不着头脑:

Exception in thread "main" java.lang.ClassNotFoundException: Failed to find data source: kafka. Please find packages at http://spark.apache.org/third-party-projects.html

我也尝试为 jar 提供 --jars 参数,但它没有帮助。我在这里想念什么?

代码:

private static void startKafkaConsumerStream() {

        Dataset<HttpPackage> ds1 = _spark
                .readStream()
                .format("kafka")
                .option("kafka.bootstrap.servers", getProperty("kafka.bootstrap.servers"))
                .option("subscribe", HTTP_FED_VO_TOPIC)
                .load() // Getting the error here
                .as(Encoders.bean(HttpPackage.class));

        ds1.foreach((ForeachFunction<HttpPackage>)  req ->System.out.print(req));

    }

而_spark定义为:

_spark = SparkSession
                .builder()
                .appName(_properties.getProperty("app.name"))
                .config("spark.master", _properties.getProperty("master"))
                .config("spark.es.nodes", _properties.getProperty("es.hosts"))
                .config("spark.es.port", _properties.getProperty("es.port"))
                .config("spark.es.index.auto.create", "true")
                .config("es.net.http.auth.user", _properties.getProperty("es.net.http.auth.user"))
                .config("es.net.http.auth.pass", _properties.getProperty("es.net.http.auth.pass"))
                .getOrCreate();

我的进口是:

import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.spark.api.java.function.ForeachFunction;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Encoders;
import org.apache.spark.sql.SparkSession;

但是,当我按照here 所述运行我的代码并且使用包选项时:

--packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.1.0

有效

【问题讨论】:

  • 我尝试添加全部详细信息,但是我无法使用更多代码提交问题。该页面不允许我!但是,请指定除上述信息之外您可能需要的信息,我会指定。
  • 您尝试使用项目访问 Kafka 的方式,因为它似乎缺少一些东西
  • 移除 kafka 依赖 = org.apache.kafka

标签: maven apache-spark apache-kafka apache-spark-sql spark-structured-streaming


【解决方案1】:

将以下依赖项添加到您的 pom.xml 文件中。

<dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql-kafka-0-10_2.11</artifactId>
        <version>2.2.0</version>
</dependency>

【讨论】:

    【解决方案2】:

    更新您的依赖项和版本。下面给定的依赖项应该可以正常工作:

        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.1.1</version>
            <scope>provided</scope>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming_2.11</artifactId>
            <version>2.1.1</version>
            <scope>provided</scope>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
            <version>2.1.1</version>
        </dependency>
    

    PS:注意前两个依赖项中提供的范围。

    【讨论】:

    • 这应该可以。您能否确保用户有权访问您在 --jars 选项中登记的 jar 文件?另外,尝试创建 fat jar 并检查它是否有效
    【解决方案3】:

    Spark 结构化流支持使用外部 kafka-0-10-sql 模块将 Apache Kafka 作为流源和接收器。

    kafka-0-10-sql 模块不适用于使用 spark-submit 提交执行的 Spark 应用程序。该模块是外部的,要使其可用,您应该将其定义为依赖项。

    除非您在 Spark 应用程序中使用 kafka-0-10-sql 特定于模块的代码,否则您不必在 pom.xml 中将模块定义为 dependency。您根本不需要模块上的编译依赖,因为没有代码使用模块的代码。您针对接口进行编码,这也是 Spark SQL 使用起来如此愉快的原因之一(即,它只需要很少的代码就可以拥有相当复杂的分布式应用程序)。

    spark-submit 但是需要--packages 命令行选项,您已报告它工作正常。

    但是,当我按照此处提到的方式运行我的代码并且使用包选项时:

    --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.1.0
    

    它与--packages 配合良好的原因是您必须告诉 Spark 基础架构在哪里可以找到kafka 格式的定义。

    这导致我们使用 Kafka 运行流式 Spark 应用程序的另一个“问题”(或要求)。您必须在 spark-sql-kafka 模块上指定 运行时依赖项。

    您可以使用--packages 命令行选项(在您spark-submit 您的 Spark 应用程序之后下载必要的 jar)或创建所谓的 uber-jar(或 fat-jar)来指定运行时依赖项。

    这就是pom.xml 发挥作用的地方(这就是人们提供pom.xml 和dependency 模块帮助的原因)。

    所以,首先,你必须在pom.xml中指定依赖。

    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-sql-kafka-0-10_2.11</artifactId>
      <version>2.2.0</version>
    </dependency>
    

    最后但并非最不重要的一点是,您必须使用Apache Maven Shade Plugin 在pom.xml 中配置一个超级jar。

    使用 Apache Maven Shade 插件,create an Uber JAR 将在 Spark 应用程序 jar 文件中包含所有 kafka 格式的“基础架构”以使其正常工作。事实上,Uber JAR 将包含所有必要的运行时依赖项,因此您可以单独使用 jar 来 spark-submit(并且没有 --packages 选项或类似选项)。

    【讨论】:

    • 尽管添加了“--packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.1.0”,我仍然得到相同的“无法找到数据源:kafka “ 错误 。你能告诉我还有什么原因吗?
    • @Bharathi 这可能有很多原因,如果不详细查看您的案例,很难在评论中回答。你能问一个单独的问题并把链接留在这里吗?谢谢。
    • 谢谢。我解决了。 scala 和 python 之间的路径似乎存在细微的语法差异。
    猜你喜欢
    • 1970-01-01
    • 2018-12-01
    • 2021-09-13
    • 2018-01-29
    • 2021-09-21
    • 2022-11-20
    • 2021-06-10
    • 2017-12-03
    • 2020-01-25
    相关资源
    最近更新 更多