【问题标题】:How to change Datatypes of records inserting into Cassandra using Foreach Spark Structure streaming如何使用 Foreach Spark 结构流更改插入 Cassandra 的记录的数据类型
【发布时间】:2019-11-21 20:05:39
【问题描述】:

我正在尝试使用 Spark Structure Streaming 和 Foreach Sink 将反序列化的 Kafka 记录插入 Data Stax Cassandra。

比如我的反序列化Data frame数据和所有的都是字符串格式的。

id   name    date
100 'test' sysdate

我使用 foreach Sink 创建了一个类,并尝试通过转换它来插入如下记录。

session.execute(
  s"""insert into ${cassandraDriver.namespace}.${cassandraDriver.brand_dub_sink} (id,name,date)
  values  ('${row.getAs[Long](0)}','${rowstring(1)}','${rowstring(2)}')"""))
  }
)

我完全按照这个项目 https://github.com/epishova/Structured-Streaming-Cassandra-Sink/blob/master/src/main/scala/cassandra_sink.scala

当插入 Cassandra 表时,如上所述将字符串“id”列数据类型转换为 Long,它没有转换。并抛出错误

“bigint 类型的“id”的 STRING 常量 (100) 无效”

卡桑德拉表;-

create table test(
id bigint,
name text,
date timestamp)

在“def Process”中将字符串数据类型转换为 Long 的任何建议。

任何替代建议也很棒。谢谢

这是代码:

import org.apache.spark.SparkContext._
import org.apache.spark.SparkConf
import org.apache.spark.sql._
import com.datastax.spark.connector._
import com.datastax.spark.connector.cql.CassandraConnector
import org.apache.spark.sql.ForeachWriter
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.functions.expr

class CassandraSinkForeach() extends ForeachWriter[org.apache.spark.sql.Row] {
  // This class implements the interface ForeachWriter, which has methods that get called 
  // whenever there is a sequence of rows generated as output

  var cassandraDriver: CassandraDriver = null;
  def open(partitionId: Long, version: Long): Boolean = {
    // open connection
    println(s"Open connection")
    true
  }

  def process(record: org.apache.spark.sql.Row) = {
    println(s"Process new $record")
    if (cassandraDriver == null) {
      cassandraDriver = new CassandraDriver();
    }
    cassandraDriver.connector.withSessionDo(session =>
      session.execute(s"""
       insert into ${cassandraDriver.namespace}.${cassandraDriver.foreachTableSink} (fx_marker, timestamp_ms, timestamp_dt)
       values('${record.getLong(0)}', '${record(1)}', '${record(2)}')""")
    )
  }

  def close(errorOrNull: Throwable): Unit = {
    // close the connection
    println(s"Close connection")
  }
}

class SparkSessionBuilder extends Serializable {
  // Build a spark session. Class is made serializable so to get access to SparkSession in a driver and executors. 
  // Note here the usage of @transient lazy val 
  def buildSparkSession: SparkSession = {
    @transient lazy val conf: SparkConf = new SparkConf()
    .setAppName("Structured Streaming from Kafka to Cassandra")
    .set("spark.cassandra.connection.host", "ec2-52-23-103-178.compute-1.amazonaws.com")
    .set("spark.sql.streaming.checkpointLocation", "checkpoint")

    @transient lazy val spark = SparkSession
    .builder()
    .config(conf)
    .getOrCreate()

    spark
  }
}

class CassandraDriver extends SparkSessionBuilder {
  // This object will be used in CassandraSinkForeach to connect to Cassandra DB from an executor.
  // It extends SparkSessionBuilder so to use the same SparkSession on each node.
  val spark = buildSparkSession

  import spark.implicits._

  val connector = CassandraConnector(spark.sparkContext.getConf)

  // Define Cassandra's table which will be used as a sink
  /* For this app I used the following table:
       CREATE TABLE fx.spark_struct_stream_sink (
       id Bigint,
       name text,
       timestamp_dt date,
       primary key (id));
  */
  val namespace = "fx"
  val foreachTableSink = "spark_struct_stream_sink"
}

object KafkaToCassandra extends SparkSessionBuilder {
  // Main body of the app. It also extends SparkSessionBuilder.
  def main(args: Array[String]) {
    val spark = buildSparkSession

    import spark.implicits._

    // Define location of Kafka brokers:
    val broker = "ec2-18-209-75-68.compute-1.amazonaws.com:9092,ec2-18-205-142-57.compute-1.amazonaws.com:9092,ec2-50-17-32-144.compute-1.amazonaws.com:9092"

    /*Here is an example massage which I get from a Kafka stream. It contains multiple jsons separated by \n 
    {"100": "test1", "01-mar-2018"}
    {"101": "test2", "02-mar-2018"}  */
    val dfraw = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", broker)
    .option("subscribe", "currency_exchange")
    .load()

    val schema = StructType(
      Seq(
        StructField("id", StringType, false),
        StructField("name", StringType, false),
StructField("date", StringType, false)

      )
    )

    val df = dfraw
    .selectExpr("CAST(value AS STRING)").as[String]
    .flatMap(_.split("\n"))

    val jsons = df.select(from_json($"value", schema) as "data").select("data.*")


    val sink = jsons
    .writeStream
    .queryName("KafkaToCassandraForeach")
    .outputMode("update")
    .foreach(new CassandraSinkForeach())
    .start()

    sink.awaitTermination()
  }
}  

我修改过的代码;-

def open(partitionId: Long, version: Long): Boolean = {
    // open connection
    println(s"in my Open connection")
    val cassandraDriver = new CassandraDriver();
    true
  }


  def process(record: Row) = {


    val optype = record(0)

    if (cassandraDriver == null) {
      val  cassandraDriver = new CassandraDriver();
    }

  if (optype == "I" || optype == "U") {

        println(s"Process insert or Update Idempotent new $record")

        cassandraDriver.connector.withSessionDo(session =>{
          val prepare_rating_brand = session.prepare(s"""insert into ${cassandraDriver.namespace}.${cassandraDriver.brand_dub_sink} (table_name,op_type,op_ts,current_ts,pos,brand_id,brand_name,brand_creation_dt,brand_modification_dt,create_date) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""")

          session.execute(prepare_rating_brand.bind(record.getAs[String](0),record.getAs[String](1),record.getAs[String](2),record.getAs[String](3),record.getAs[String](4),record.getAs[BigInt](5),record.getAs[String](6),record.getAs[String](7),record.getAs[String](8),record.getAs[String](9))
          )

        })
      }  else if (optype == "D") {

        println(s"Process delete new $record")
        cassandraDriver.connector.withSessionDo(session =>
          session.execute(s"""DELETE FROM ${cassandraDriver.namespace}.${cassandraDriver.brand_dub_sink} WHERE brand_id = ${record.getAs[Long](5)}"""))

      } else if (optype == "T") {
        println(s"Process Truncate new $record")
        cassandraDriver.connector.withSessionDo(session =>
          session.execute(s"""Truncate table  ${cassandraDriver.namespace}.${cassandraDriver.plan_rating_archive_dub_sink}"""))

      }
    }

  def close(errorOrNull: Throwable): Unit = {
    // close the connection
    println(s"Close connection")
  }


}

【问题讨论】:

  • 如果您使用的是 Spark 2.4.x,那么您可以改用 foreachBatch - 它比 foreach 简单得多。这是一个例子:github.com/alexott/spark-cassandra-oss-demos/blob/master/src/…
  • 谢谢,@Alex Ott。我正在使用 Spark 2.2 和 datastax 6.0。
  • 如果您使用 DataStax 发行版,您可以简单地使用 writeStream,而不使用此处理器:github.com/alexott/dse-java-playground/blob/master/src/main/… - 在同一文件夹中的另一个文件中,有一个 Kafka 示例
  • 好的,@AlexOtt,考虑这种情况,我有 10 个查询写入 datastax,一个查询在 spark 中终止。我有流式查询管理器、检查点和异常处理一切。您是否指导我如何在不影响会话的情况下重新启动该特定查询?因为我已经处理了 queryawaittermination()。
  • 嗯,我没有完全理解。如果可能,通常会自动重试查询:docs.datastax.com/en/developer/java-driver/3.7/manual/retries 但您可能需要专门将其标记为幂等。但正如我所提到的,如果您使用的是 DataStax Enterprise - 您不需要自己编写所有这些代码 - 您只需要使用 writeStream - 然后 DSE 连接器将自动处理所有内容

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


【解决方案1】:

您的错误是您将 id 字段的值指定为 '${row.getAs[Long](0)}' - 您在其周围添加了单引号,因此它被视为字符串,而不是 long/bigint - 只需删除单个围绕这个值引用:${row.getAs[Long](0)}...

另外,出于性能原因,最好将 cassandra 驱动程序的实例化移动到 open 方法中,并使用准备好的语句,如下所示:

  var cassandraDriver: CassandraDriver = null;
  var preparedStatement: PreparedStatement = null;
  def open(partitionId: Long, version: Long): Boolean = {
    // open connection
    println(s"Open connection")
    cassandraDriver = new CassandraDriver();
    preparedStatement = cassandraDriver.connector.withSessionDo(session =>
      session.prepare(s"""
       insert into ${cassandraDriver.namespace}.${cassandraDriver.foreachTableSink} 
      (fx_marker, timestamp_ms, timestamp_dt) values(?, ?, ?)""")
    true
  }

  def process(record: org.apache.spark.sql.Row) = {
    println(s"Process new $record")
    cassandraDriver.connector.withSessionDo(session =>
      session.execute(preparedStatement.bind(${record.getLong(0)}, 
           ${record(1)}, ${record(2)}))
    )
  }

它会更高效,而且您不需要自己执行值的引用。

【讨论】:

  • 嗨,Alex,还有一个疑问。我正在使用另外一个参数 op_type = I(Insert) 或 U(Update) 接收从金门到 Kafka 的所有数据。因此,我将根据必须在 Cassandra 中插入的数据来使用数据 [插入或更新]。所以我必须在 [def Process] 中处理 Prepared 语句。这可能吗
  • 所以只需准备 2 个不同的查询,并根据标志选择它们。
  • 谢谢@Alex Ott。 datastax-oss.atlassian.net/browse/…。我从 Data stax 中读取它。所以 foreach 内部有准备语句?
  • 啊,是的 - 忘记了...您可以将此字段标记为@transient - 请参阅waitingforcode.com/apache-spark/serialization-issues-part-2/…
  • 嗨,亚历克斯,还有一个疑问,当我尝试在列上更新时,它的字符串类似于“Bravo Daytime-D (MF 8 AM-3 PM)”,我收到了类似 com 的错误.datastax.driver.core.exceptions.SyntaxError:第 60:24 行在输入“Bravo”处没有可行的替代方案(... =63904,sales_unit_name = ''[Bravo]...)如何动态使用转义字符。 ? @亚历克斯·奥特
猜你喜欢
  • 2018-10-06
  • 2016-07-18
  • 1970-01-01
  • 2018-11-18
  • 2019-11-29
  • 2023-03-21
  • 2020-06-16
  • 2021-03-16
  • 2017-10-05
相关资源
最近更新 更多