【问题标题】:Spark Streaming: NullPointerException inside foreachPartitionSpark Streaming:foreachPartition 内的 NullPointerException
【发布时间】:2016-02-05 12:00:12
【问题描述】:

我有一个火花流作业,它从 Kafka 读取数据,并在再次写入 Postrges 之前与 Postgres 中的现有表进行一些比较。这就是它的样子:

val message = KafkaUtils.createStream(...).map(_._2)

message.foreachRDD( rdd => {

  if (!rdd.isEmpty){
    val kafkaDF = sqlContext.read.json(rdd)
    println("First")

    kafkaDF.foreachPartition(
      i =>{
        val jdbcDF = sqlContext.read.format("jdbc").options(
          Map("url" -> "jdbc:postgresql://...",
            "dbtable" -> "table", "user" -> "user", "password" -> "pwd" )).load()

        createConnection()
        i.foreach(
          row =>{
            println("Second")
            connection.sendToTable()
          }
        )
        closeConnection()
      }
    )

这段代码在 val jbdcDF = ... 行给了我 NullPointerException

我做错了什么?此外,我的日志"First" 有效,但"Second" 没有出现在日志中的任何位置。我用kafkaDF.collect().foreach(...) 尝试了整个代码,它运行良好,但性能很差。我希望用foreachPartition 替换它。

谢谢

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    目前尚不清楚createConnectioncloseConnectionconnection.sendToTable 内部是否存在任何问题,但根本问题是尝试嵌套动作/转换。 Spark 不支持它,Spark Streaming 也不例外。

    这意味着嵌套的DataFrame 初始化 (val jdbcDF = sqlContext.read.format ...) 根本无法工作,应该被删除。如果您将其用作参考,则应在与kafkaDF 相同的级别创建它并使用标准转换(unionAlljoin、...)进行参考。

    如果由于某种原因它不是一个可接受的解决方案,您可以在 forEachPartition 内创建纯 JDBC 连接并在 PostgreSQL 表上进行操作(我猜这是您在 sendToTable 内已经做的事情)。

    【讨论】:

    • 我在没有 jdbcDF 初始化的情况下尝试过,它正在工作,所以绝对不是任何连接操作的问题
    • 我不能使用像 join 这样的运算符,因为我想遍历 kafkaDF 中的每一行并与 jdbcDF 进行一些比较。
    • 在这里您无能为力。您可以在 Spark 中以一种或另一种方式编码您的逻辑(推荐方法)或创建标准 JDBC 连接(不是 Spark DataFrame)并直接访问 PostgreSQL 中的数据。
    • 请检查这个问题。 stackoverflow.com/questions/32458109/…
    • 他们正在讨论解决方案。但是我在创建与 kafkaDF 相同级别的 jdbcDF 后尝试广播它,仍然给我 NPE
    【解决方案2】:

    正如@zero323 正确指出的那样,你不能广播你的jdbc 连接,你也不能创建嵌套的RDD。 Spark 根本不支持在现有闭包(即 foreachPartition)中使用 sparkContext 或 sqlContext,因此会出现空指针异常。

    有效解决此问题的唯一方法是在 foreachPartition 中创建一个 JDBC 连接并直接在其上执行 SQL 以执行您想要的任何操作,然后使用相同的连接写回记录。

    关于您的第二个已编辑问题:

    变化:

    kafkaDF.foreachPartition(..)
    

    kafkaDF.repartition(numPartition).foreachPartition(..)
    

    其中 numPartition 是所需的分区数。这将增加分区的数量。如果您有多个执行器(每个执行器有多个任务),这些将并行运行。

    【讨论】:

    猜你喜欢
    • 2014-12-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-05-19
    • 2021-02-01
    • 2017-03-17
    • 2021-03-14
    相关资源
    最近更新 更多