【问题标题】:How can I get two different cassandra clusters in my spark structured streaming?如何在我的 spark 结构化流中获得两个不同的 cassandra 集群?
【发布时间】:2021-09-20 01:46:15
【问题描述】:

我有两个 cassandra 集群。如何在同一个 SparkSession 中设置主机、密码和用户?以及如何将它们与 CassandraConnector 一起使用?

我试过了:

val cassandraCon: CassandraConnector = CassandraConnector(conf)

  val ks = "monitore"
  val ttableName = "validate_structure"

  def getIndex(): ResultSet = {

    val table = ks + "." + ttableName
    val query = s"""select *
                   |from ${table}""".stripMargin
    println(query)
    cassandraCon.withSessionDo(s => {
      s.execute(query)
    })
  }

但是,问题在于,这仅在 cassandra 集群与 spark 位于同一主机上时才有效。我也尝试创建一个目录,但我找不到使用 session.execute 而不是 spark.sql 发出请求的方法。

有人可以帮助我吗?我使用 Spark Structured Streaming 3.1.2 并使用 cassandra 连接来丰富我的数据。

【问题讨论】:

  • 你想达到什么目的?对于 select * 你不需要 withSessionDo
  • @AlexOtt 在 select 中我将放置一个 where 子句,该子句将从数据集中接收数据作为参数。
  • 你不需要这样做,这真的是一个糟糕的做法 - 我明天会回答

标签: apache-spark spark-structured-streaming spark-cassandra-connector


【解决方案1】:

第一个问题是.withSessionDo仅在驱动程序的上下文中运行,而不是在执行程序的上下文中,所以它不会被分发。

你需要使用:

  • 来自 RDD API 的 .joinWithCassandraTable function(也有它的左连接版本)
  • 或使用所谓的DirectJoin(请参阅blog post of its author 中的详细信息),当 Spark Cassandra 连接器 (SCC) 检测到联接的一侧在 Cassandra 中时,并将其转换为对单个分区的查询。不幸的是,Spark 3.1 在 SCC 中破坏了当前版本的 DirectJoin(请参阅 JIRA),因此您可能需要使用 RDD API,直到它得到修复。

我有 detailed blog post 关于如何在 Cassandra 中对数据执行高效联接的信息。

关于两个集群 - 完全可以为单个读/写操作指定连接详细信息,只需指定 .option("spark.cassandra.connection.host", "host-or-ip")。 (Russell Spitzer,SCC 的主要开发人员,拥有blog posts on how to connect to multiple clusters)。当您使用目录 API 时也可以这样做 - 只需将连接属性名称(例如,spark.cassandra.connection.host)附加到特定目录名称(请参阅docs

【讨论】:

  • 感谢您的帮助!
猜你喜欢
  • 1970-01-01
  • 2017-08-25
  • 2017-09-16
  • 1970-01-01
  • 2019-11-12
  • 2018-12-16
  • 2018-10-06
  • 2019-07-05
  • 1970-01-01
相关资源
最近更新 更多