You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.23 14:17:48