【问题标题】:How to split text in each row when getting the data from kafka topic?从kafka主题获取数据时如何拆分每行中的文本?
【发布时间】:2021-02-05 17:32:13
【问题描述】:

我从一个 kafka 主题中检索数据。将数据转换为具有 10 列的数据框后,我选择其中一列。我想在每一行中拆分字符串,以便将单词转换为它们的发音。 evrything 似乎没问题,但我唯一的问题是我不能连续运行 split 方法。?

这是我的代码

   val df = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "rocket-01.srvs.cloudkafka.com:9094")
            .option("subscribe", "k-news") 
            .option("startingOffsets", "latest")
            .option("kafka.security.protocol","SASL_SSL")
            .option("kafka.sasl.mechanism", "SCRAM-SHA-256")
            .option("kafka.sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=" " password="";").load()

val newsStringDF = df.selectExpr("CAST(value AS STRING)")

val newsSchema = StructType( Array(
     StructField("_c0",StringType,true),
     StructField("id",StringType,true),
     StructField("title",StringType,true),
     StructField("publication",StringType,true),
     StructField("author",StringType,true),
     StructField("date",StringType,true),
     StructField("year",StringType,true),
     StructField("month",StringType,true),
     StructField("url",StringType,true),
     StructField("content",StringType,true)))


val lines = Source.fromFile("/Project/symbols/cmudict.dict").getLines()
val res = lines.map { line =>line.split(" ", 2) match { case Array(a, b) => a -> b }}.toMap


val newsDF = newsStringDF.select(from_json(col("value"),
newsSchema).as("data")).select("data.*")
val titleColumn = newsDF.select("title").as[String].foreach( message => 
message.split(" ").toList.map { s =>if(res.get(s).isDefined) {res(s)} else 
{s}}.mkString(" ")).toDF("title") 

val streaming = titleColumn.writeStream.format("console").outputMode("append").trigger(Trigger.ProcessingTime("10 seconds")).start().awaitTermination()

输出应该是这样的: 我从我的 kafka-topic 获取消息(来自标题列的字符串),然后用另一个字符串替换它。 “我喜欢足球” --> “我喜欢 ˈfo͝otˌbôl” 并使用 .writeStream 将消息写入控制台

提前致谢

【问题讨论】:

  • 你能用最少的代码来重现这个吗? (包括模式定义、样本输入数据、变量res 的定义)。理想情况下,还应提供预期输出的示例。
  • 最好能得到一个样本输入和输出数据集和几行代码sn-p,你必须清楚地了解问题。
  • 现在我更新了我的问题并试图解释输出应该是什么样子。感谢您提前提供任何帮助。
  • 你能告诉我你的 spark-submit 命令吗?我正在处理类似的问题,但我不确定如何配置 spark-submit

标签: scala apache-spark apache-spark-sql spark-streaming apache-kafka-streams


【解决方案1】:

您必须记住,Dataframe 实际上只是 Dataset[Row]。当使用简单类型(甚至一些自定义案例类)时,很容易在有类型的Dataset 和无类型的Dataframe 之间进行转换。在这种特殊情况下,它应该很简单:

val titleColumn = newsDF
  .select("title")
  .as[String]
  .map( message => message.split(" ").toList.map { s =>
    if(res.get(s).isDefined) {res(s)} else {s}}.mkString(" ")
  )
  .toDF("title")

额外的行.as[String] 将您的Dataframe 转换为Dataset[String].toDF("title") 将其转换回来。这允许您的地图在Strings 而不是Rows 上运行。由于这是一个转换,您还需要在Dataset 上使用.map 而不是.foreach

另一种选择是从Row 中检索标题String 并对其进行拆分:

message.getAs[String]("title").split(" ")

我不确定哪种方法更有效。这是一个经过测试的样本:

val df = Seq("aa 11", "bb 22", "cc 33", "dd 44").toDF("title")

val res = Map("aa" -> "AA", "bb" -> "BB", "cc" -> "CC", "dd" -> "DD")

df.as[String].map(m => {
  m.split(" ").toList.map(s => {
    if(res.get(s).isDefined) res(s) else s    
  }).mkString(" ")
}).toDF("title").show()

结果:

+-----+
|title|
+-----+
|AA 11|
|BB 22|
|CC 33|
|DD 44|
+-----+

【讨论】:

  • 不幸的是,这两种方法都不适合我。我尝试了第一种方法,我得到一个错误:重载方法值 foreach with alternatives: (func: org.apache.spark.api.java.function.ForeachFunction[String])Unit (f: String => Unit)单元。当我尝试其他方法时,我得到了几乎相同的错误:
  • 哎呀!我第一次没有注意到.foreach。您需要使用.map,因为这是一个转换。我认为该错误与您的匿名函数返回 List[String] 而不是 Unit 的事实有关。
  • 查看编辑,我在其中包含了一个经过测试的示例代码块。
  • 是的,一旦我用地图换了 foreach,它就起作用了。所以这部分代码不再是问题:),但是当我尝试将 writeStream 写入控制台时出现错误。 org.apache.spark.SparkException:在 org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:416) 上无法序列化的任务。试图搜索可以导致此错误的原因,但仍然没有找到
  • 当您在匿名函数内部引用无法序列化的变量时会发生该错误。 .writeStream 方法是您的操作,这就是导致转换代码实际执行的原因。该错误可能在此之前发生。通常,您会在匿名函数内部看到这个不可序列化的变量,这些函数会分发到集群(mapflatMapforeach 等针对数据帧)。很多时候,它是一个代表数据库连接之类的实例,不能移动到其他节点。
猜你喜欢
  • 2016-05-28
  • 1970-01-01
  • 2021-10-07
  • 1970-01-01
  • 2021-07-03
  • 2019-06-12
  • 2019-01-04
  • 2018-09-20
  • 2023-02-01
相关资源
最近更新 更多