【问题标题】:How to use writeStream to pass Spark stream to a kafka topic如何使用 writeStream 将 Spark 流传递给 kafka 主题
【发布时间】:2020-03-08 21:34:18
【问题描述】:

我正在使用提供流的 twitter 流功能。我需要使用 Spark writeStream 函数,例如:writeStream function link

// Write key-value data from a DataFrame to a specific Kafka topic specified in an option
val ds = df
  .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
  .option("topic", "topic1")
  .start()

“df”需要是流数据集/数据帧。如果 df 是一个普通的 DataFrame,它会报错,显示 'writeStream' 只能在流数据集/DataFrame 上调用;

我已经完成了: 1. 从推特获取信息流 2.过滤和映射得到每个twitt的标签(Positive, Negative, Natural)

最后一步是 groupBy 标记和计数,并将其传递给 Kafka。

你们知道如何将 Dstream 转换为流式 Dataset/DataFrame 吗?

已编辑:ForeachRDD 函数确实将 Dstream 更改为普通 DataFrame。 但是“writeStream”只能在流式传输时调用 数据集/数据框。 (上面提供了writeStream链接)

org.apache.spark.sql.AnalysisException: 'writeStream' 只能在流数据集/DataFrame 上调用;

【问题讨论】:

标签: apache-kafka spark-streaming spark-structured-streaming


【解决方案1】:

如何将 Dstream 转换为流式 Dataset/DataFrame?

DStream 是一系列 RDD 的抽象。

流式Dataset 是一系列Datasets 的“抽象”(我使用引号,因为流式和批处理Datasets 之间的区别是isStreamingDataset 的属性)。

可以将DStream 转换为流式Dataset 以保持DStream 的行为。

我认为你并不是真的想要它。

您只需要使用DStream 获取推文并将它们保存到 Kafka 主题(并且您认为您需要结构化流)。我认为您只需要 Spark SQL(结构化流的底层引擎)。

然后伪代码如下(抱歉,自从我使用老式的 Spark Streaming 以来已经有一段时间了):

val spark: SparkSession = ...
val tweets = DStream...
tweets.foreachRDD { rdd =>
  import spark.implicits._
  rdd.toDF.write.format("kafka")...
}

【讨论】:

  • 我试过这个方法,但它给了我错误。最后,我使用 KafkaProducer 向我的主题写入消息。
  • @DDJin 有什么错误?如果您不介意,我可以帮助您:)
猜你喜欢
  • 2022-07-02
  • 2019-05-13
  • 2018-03-06
  • 2019-07-21
  • 2018-01-10
  • 2016-10-15
  • 2019-04-23
  • 1970-01-01
相关资源
最近更新 更多