Spark Structured Streaming代码报错:无法找到Trigger.once
你的报错 error: not found: value Trigger 背后藏着两个核心问题,我们一步步拆解:
1. 调用时机完全错误:trigger 属于 DataStreamWriter,而非已启动的 StreamingQuery
你已经通过 .start() 得到了 populationStream(这是一个已启动的流查询对象),但 trigger() 方法只能在构建流写入器的阶段使用——也就是在调用 .start() 之前设置。
你原代码里已经设置过 trigger(Trigger.ProcessingTime(30.seconds)),之后试图在已运行的查询上再次调用 trigger,这是不符合Spark API设计的错误用法。
2. Spark版本不支持 Trigger.once()
从你提供的依赖包(比如 spark-sql-kafka-0-10_2.11-2.1.1.jar)能看出,你用的是Spark 2.1.1,而 Trigger.once() 是Spark 2.2.0才引入的新特性。在2.1.1版本里,这个API根本不存在,哪怕你导入了正确的包,也会找不到这个方法。
针对性修正方案
如果你需要「执行一次就停止」的效果(Spark 2.1.1适配版)
在Spark 2.1.1中没有Trigger.once(),可以用以下两种方式实现类似效果:
- 方式一:用
ProcessingTime(0)触发立即处理,完成后手动停止
val populationStream = population .writeStream .format("console") .trigger(Trigger.ProcessingTime(0.seconds)) // 立即执行一次处理 .outputMode(OutputMode.Complete) .queryName("textStream") .start() populationStream.awaitTermination() // 等待处理完成 populationStream.stop() // 主动停止查询
- 方式二:直接用批处理读取(如果你的CSV数据源是静态文件,其实没必要用流处理)
val people = spark.read.schema(schema).csv("/data/pncdw/scratch/test/*.csv").as[Person] val population = people.groupBy('city).agg(count('city) as "population") population.show()
如果一定要用 Trigger.once()
你需要把Spark版本升级到2.2.0或更高,同时确保所有依赖包的版本和Spark版本一致(比如把spark-sql-kafka-0-10_2.11-2.1.1.jar换成对应2.2.0版本的包),然后在流写入器构建阶段正确设置:
val populationStream = population .writeStream .format("console") .trigger(Trigger.once()) // 仅Spark 2.2+可用 .outputMode(OutputMode.Complete) .queryName("textStream") .start()
另外提个小细节:你代码里重复定义了schema变量(先手动创建StructType,又用Encoders.product[Person].schema覆盖),虽然不影响运行,但建议保持变量定义的一致性,避免后期混淆。
内容的提问来源于stack exchange,提问作者Prasad Ligade

