Spark流处理报错:value select不是DataStreamReader的成员
问题解决方法
你的报错核心原因是:kafkaDF是DataStreamReader类型的对象,而select是Spark DataFrame/Dataset才有的方法,直接在DataStreamReader上调用肯定会报错。
解决步骤很简单,先把DataStreamReader转换成Streaming DataFrame:
- 先调用
load()方法:
val kafkaStreamDF = kafkaDF.load()
- 再对转换后的Streaming DataFrame执行select操作:
val activationDF = kafkaStreamDF.select(from_json($"value".cast("string"), activationSchema).alias("activation"))
也可以把两步合并成一行:
val activationDF = kafkaDF.load().select(from_json($"value".cast("string"), activationSchema).alias("activation"))
补充说明:在Spark结构化流处理中,spark.readStream.format("kafka")这类代码返回的是DataStreamReader,它只是一个读取流数据的配置器,必须调用load()才能真正获取到可以进行数据转换的Streaming DataFrame,之后才能用select、filter这些DataFrame方法。
内容的提问来源于stack exchange,提问作者Melih
相关产品推荐
相关产品推荐

