【问题标题】:java.io.IOException: Failed to write statements to batch_layer.test. The latest exception was Key may not be emptyjava.io.IOException:无法将语句写入 batch_layer.test。最新的例外是 Key may not be empty
【发布时间】:2021-07-30 18:12:06
【问题描述】:

我正在尝试计算文本中的单词数并将结果保存到 Cassandra 数据库。 Producer 从文件中读取数据并发送给kafka。 Consumer 使用 Spark Streaming 读取和处理日期,然后将计算结果发送到表中。

我的制作人是这样的:

object ProducerPlayground extends App {

  val topicName = "test"
  private def createProducer: Properties = {
    val producerProperties = new Properties()
    producerProperties.setProperty(
      ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
      "localhost:9092"
    )
    producerProperties.setProperty(
      ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
      classOf[IntegerSerializer].getName
    )
    producerProperties.setProperty(
      ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      classOf[StringSerializer].getName
    )
    producerProperties
  }

  val producer = new KafkaProducer[Int, String](createProducer)

  val source = Source.fromFile("G:\\text.txt", "UTF-8")

  val lines = source.getLines()

  var key = 0
  for (line <- lines) {
    producer.send(new ProducerRecord[Int, String](topicName, key, line))
    key += 1
  }
  source.close()
  producer.flush()

}

消费者看起来像这样:

object BatchLayer {
  def main(args: Array[String]) {

    val brokers = "localhost:9092"
    val topics = "test"
    val groupId = "groupId-1"

    val sparkConf = new SparkConf()
      .setAppName("BatchLayer")
      .setMaster("local[*]")
    val ssc = new StreamingContext(sparkConf, Seconds(3))
    val sc = ssc.sparkContext
    sc.setLogLevel("OFF")

    val topicsSet = topics.split(",").toSet
    val kafkaParams = Map[String, Object](
      ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG -> brokers,
      ConsumerConfig.GROUP_ID_CONFIG -> groupId,
      ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG -> classOf[StringDeserializer],
      ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG -> classOf[StringDeserializer],
      ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG -> "false"
    )
    val stream =
      KafkaUtils.createDirectStream[String, String](
        ssc,
        LocationStrategies.PreferConsistent,
        ConsumerStrategies.Subscribe[String, String](topicsSet, kafkaParams)
      )

   
    val cass = CassandraConnector(sparkConf)

    cass.withSessionDo { session =>
      session.execute(
        s"CREATE KEYSPACE IF NOT EXISTS batch_layer WITH REPLICATION = {'class': 'SimpleStrategy', 'replication_factor': 1 }"
      )
      session.execute(s"CREATE TABLE IF NOT EXISTS batch_layer.test (key VARCHAR PRIMARY KEY, value INT)")
      session.execute(s"TRUNCATE batch_layer.test")
    }

    stream
      .map(v => v.value())
      .flatMap(x => x.split(" "))
      .filter(x => !x.contains(Array('\n', '\t')))
      .map(x => (x, 1))
      .reduceByKey(_ + _)
      .saveToCassandra("batch_layer", "test", SomeColumns("key", "value"))

    ssc.start()
    ssc.awaitTermination()
  }

}

启动生产者后,程序停止工作并出现此错误。我做错了什么?

【问题讨论】:

  • 添加您遇到的错误
  • @AlexOtt,java.io.IOException: Failed to write statements to batch_layer.test. The latest exception was Key may not be empty。此错误出现两次,程序以代码 1 终止
  • 看起来您的 key 具有空值

标签: apache-kafka spark-streaming cassandra-3.0 spark-cassandra-connector spark-streaming-kafka


【解决方案1】:

在 2021 年使用传统流媒体几乎没有意义 - 使用起来非常麻烦,而且您还需要跟踪 Kafka 等的偏移量。最好使用 Structured Streaming instead - 它会通过检查点跟踪您的偏移量,您将使用高级数据集 API 等。

在您的情况下,代码可能如下所示(未测试,但从 this working example 采用):

val streamingInputDF = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9092")
  .option("subscribe", "test")
  .load()

val wordsCountsDF = streamingInputDF.selectExpr("CAST(value AS STRING) as value")
  .selectExpr("split(value, '\\w+', -1) as words")
  .selectExpr("explode(words) as word")
  .filter("word != ''")
  .groupBy($"word")
  .count()
  .select($"word", $"count")

// create table ...

val query = wordsCountsDF.writeStream
   .outputMode(OutputMode.Update)
   .format("org.apache.spark.sql.cassandra")
   .option("checkpointLocation", "path_to_checkpoint)
   .option("keyspace", "test")
   .option("table", "<table_name>")
   .start()

query.awaitTermination()

附:在您的示例中,最可能的错误是您尝试直接在 DStream 上使用 .saveToCassandra - 它不能以这种方式工作。

【讨论】:

  • 另外值得指出的是,生产者是不必要的,因为 Spark 可以用来读取输入文件。然后可以将该文件数据直接写入 Cassandra
猜你喜欢
  • 2019-12-06
  • 1970-01-01
  • 2014-01-08
  • 1970-01-01
  • 1970-01-01
  • 2015-12-17
  • 2010-10-25
  • 2019-04-07
  • 1970-01-01
相关资源
最近更新 更多