Spark Streaming中foreachBatch未执行,run_command函数未运行问题排查求助
Spark Streaming中foreachBatch未执行,run_command函数未运行问题排查求助
大家好,我碰到了一个Structured Streaming的问题:下面这段代码里的run_command函数似乎根本没有被执行,连函数里的打印语句都没输出,想请各位帮忙看看问题出在哪?
我的代码如下:
df.writeStream \ .foreachBatch(run_command) \ .format("delta") \ .outputMode("append") \ .option("checkpointLocation", "/tmp/delta/events/_checkpoints/") \ .partitionBy("machineid") \ .toTable("mytimeseriesdata") n = 100 count = 0 def run_command(batchDF, epoch_id): global count print("something") count += 1 if count % n == 0: spark.sql("OPTIMIZE mytimeseriesdata ZORDER BY (timestamp)") print("Optimizing " + count)
我原本的预期是,每处理100个批次就自动执行一次Delta表的OPTIMIZE优化,但现在不仅优化没做,连函数里最基础的print("something")都没看到任何输出,感觉整个函数都没被触发。有没有大佬能帮忙分析下原因?
备注:内容来源于stack exchange,提问作者Mathias Rönnlund
相关产品推荐
相关产品推荐

