Spark Streaming中foreach及foreachBatch使用报错与自定义操作实现问题
问题原因
foreachBatch()只打印一次:该API要求传入可调用的函数对象,你直接填写print("do i get printed every batch?")会在代码初始化阶段就立即执行print语句,把返回的None传给foreachBatch(),所以只有启动时打印一次,后续批次不会触发执行。foreach()报错:foreach()是针对单条数据做处理的API,要求传入的参数要么是实现了process方法的类实例,要么是可直接调用的函数对象,你还是传入了print语句执行后的返回值None,自然会报缺少process方法的错误,且foreach()本身也不符合你每批次执行一次自定义操作的需求。
解决方案
使用foreachBatch实现需求,按如下逻辑调整代码:
- 先定义批次处理函数,函数固定接收两个参数:当前批次的DataFrame、批次ID
- 在函数内先执行你的自定义操作(打印、其他业务逻辑都可在此添加),再完成BigQuery的写入
- 把定义好的函数名传给
foreachBatch()即可
# 自定义批次处理函数 def process_batch(batch_df, batch_id): # 每批次执行的自定义操作写在这里 print(f"do i get printed every batch? 当前批次ID:{batch_id}") # 批次数据写入BigQuery的逻辑迁移到函数内部 batch_df.write.format("bigquery").mode("append") \ .option("temporaryGcsBucket", path1) \ .option("table", table_kafka) \ .save() # 流处理启动逻辑 batch_job = df_alarmsFromKafka.writeStream \ .trigger(processingTime='120 seconds') \ .foreachBatch(process_batch) \ .outputMode("append") \ .option("checkpointLocation", path2) \ .start() batch_job.awaitTermination()
注意:如果是集群部署模式,自定义操作里的标准输出会打印在对应executor的日志中,local模式下可以正常在Jupyter输出单元看到打印内容。
内容的提问来源于stack exchange,提问作者nonoDa
相关产品推荐
相关产品推荐

