【问题标题】:How to know in scala runtime the different join types spark如何在 scala 运行时知道不同的连接类型火花
【发布时间】:2019-01-04 08:46:25
【问题描述】:

我想根据可用的 Spark 连接类型白名单测试用户输入。

有没有办法通过内置的 spark 了解不同的连接类型?

例如,我想根据Seq Seq("inner", "cross", "outer", "full", "fullouter", "left", "leftouter", "right", "rightouter", "leftsemi", "leftanti") 验证用户的输入

(哪些是 Spark 中可用的所有连接类型)无需像我刚刚所做的那样对其进行硬编码。

【问题讨论】:

  • 您能举例说明您需要什么以及您期望的类型吗?
  • 我想知道是否有办法获得像Seq("inner", "cross", "outer", "full", "fullouter", "left", "leftouter", "right", "rightouter", "leftsemi", "leftanti") 这样的东西,例如我可以验证用户的输入:)
  • 您需要在运行时验证连接类型并选择最接近的连接类型吗?或者您需要在运行时传递连接类型?
  • 例如,如果用户的输入不是不同的 joinType 之一,我想抛出异常。但我想避免像我刚刚所做的那样对连接类型进行硬编码
  • 检查我的答案

标签: scala apache-spark join apache-spark-sql


【解决方案1】:

我改编了这个问题here 的答案。您还可以在 Json 文件中添加 joinTypes 以在 runtume 中读取。您可以查看此答案以了解 json 对象处理JsonParsing

更新1:我更新答案以遵循Spark文档方式JoinType

import org.apache.spark._
import org.apache.spark.sql._
import org.apache.spark.sql.expressions._
import org.apache.spark.sql.functions._


object SparkSandbox extends App {

  case class Row(id: Int, value: String)

  private[this] implicit val spark = SparkSession.builder().master("local[*]").getOrCreate()

  import spark.implicits._

  spark.sparkContext.setLogLevel("ERROR")

  val r1 = Seq(Row(1, "A1"), Row(2, "A2"), Row(3, "A3"), Row(4, "A4")).toDS()
  val r2 = Seq(Row(3, "A3"), Row(4, "A4"), Row(4, "A4_1"), Row(5, "A5"), Row(6, "A6")).toDS()
  val validUserJoinType = "inner"
  val inValiedUserJoinType = "nothing"

  val joinTypes = Seq("inner", "outer", "full", "full_outer", "left", "left_outer", "right", "right_outer", "left_semi", "left_anti")

  inValiedUserJoinType match {
    case x => if (joinTypes.contains(x)) {
      println("do some logic")
      joinTypes foreach { joinType =>
        println(s"${joinType.toUpperCase()} JOIN")
        r1.join(right = r2, usingColumns = Seq("id"), joinType = joinType).orderBy("id").show()
      }
    }
    case _ =>
  val supported = Seq(
    "inner",
    "outer", "full", "fullouter", "full_outer",
    "leftouter", "left", "left_outer",
    "rightouter", "right", "right_outer",
    "leftsemi", "left_semi",
    "leftanti", "left_anti",
    "cross")

  throw new IllegalArgumentException(s"Unsupported join type '$inValiedUserJoinType'. " +
  "Supported join types include: " + supported.mkString("'", "', '", "'") + ".")
  }

}

【讨论】:

【解决方案2】:

抱歉,如果没有对 Spark 项目本身的 PR,这是不可能的。连接类型在 JoinType 内联定义。有扩展 JoinType 的类,但命名约定与 case 语句中使用的字符串不同。所以恐怕你运气不好。

https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/joinTypes.scala

【讨论】:

  • 这是我的确切问题,我并不特别需要它,更多的是为了兴趣
  • 感谢您为避免硬编码所做的努力!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-12-09
  • 1970-01-01
相关资源
最近更新 更多