【发布时间】: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 上调用;
【问题讨论】:
-
@cricket_007 感谢您的信息。我以前做过。我将问题更新为更具体。
-
啊,显然不可能... stackoverflow.com/questions/49559007/… 但是你可以写成批处理而不是spark.apache.org/docs/latest/…
-
什么是“推特流功能”?你能显示执行此操作的代码吗?
标签: apache-kafka spark-streaming spark-structured-streaming