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

Azure Databricks中如何从流式Delta表逐行获取数据?

在Azure Databricks中从流式Delta表逐行获取数据的解决方案

问题原因

你在spark-shell/cmd中能用println看到输出,但在Databricks中无法生效,核心原因是:

  • Databricks流作业的executor进程标准输出(stdout)不会直接转发到笔记本界面,println的内容仅会写入executor的日志文件,无法在笔记本输出区域显示。
  • 本地spark-shell的执行模式与Databricks集群模式存在差异,直接用println并非Databricks流式处理的标准输出/调试方式。

可行解决方案

1. 使用foreachBatch进行逐行处理(推荐用于调试或批量逐行场景)

foreachBatch允许你在每个微批次中获取完整的DataFrame,再遍历每一行处理。你可以将结果写入其他Delta表、输出到集群日志,或临时展示(仅调试用)。

示例代码:

val process_deltatable = read_deltatable.writeStream.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
  // 遍历批次内的每一行
  batchDF.foreach { row =>
    val telemetry = row.getString(0)
    // 方式1:写入集群日志(可在集群日志页面查看)
    org.apache.log4j.Logger.getLogger("DeltaStreamProcessor").info(s"Telemetry data: $telemetry")
    // 方式2:小批量调试时,可收集数据后用display展示(大数据量禁用,避免内存溢出)
  }
  // 可选:将处理后的数据写入目标Delta表
  // batchDF.write.format("delta").mode("append").saveAsTable("processed_telemetry")
}

val xyz = process_deltatable.start()

2. 改进ForeachWriter,用日志或外部存储输出

如果必须使用ForeachWriter(比如需要严格的逐行流处理),不要用println,改用日志框架输出到集群日志,或把数据写入外部系统(如ADLS、Kafka):

示例代码(日志输出):

import org.apache.log4j.Logger

val process_deltatable = read_deltatable.writeStream.foreach(new ForeachWriter[Row] {
  private val logger = Logger.getLogger(classOf[ForeachWriter[Row]])
  
  def process(value: Row): Unit = {
    val telemetry = value.getString(0)
    logger.info(s"Processing telemetry: $telemetry")
  }
  
  def open(partitionId: Long, epochId: Long): Boolean = true
  
  def close(errorOrNull: Throwable): Unit = {}
})

val xyz = process_deltatable.start()

你可以通过Databricks集群的日志页面(集群详情→日志→executor日志)查看输出内容。

3. 控制台输出(仅调试用,生产环境禁用)

如果只是临时调试,可直接用format("console"),Databricks会将流输出展示在笔记本中:

val process_deltatable = read_deltatable.writeStream
  .format("console")
  .option("truncate", false) // 完整显示内容
  .start()

注意:该方式仅适用于小数据量调试,大数据量下会严重影响集群性能。

额外注意事项

  • 确保read_deltatable是正确的流式Delta表读取对象,需用spark.readStream.format("delta").load("<delta-path>")创建。
  • 流作业启动后,可通过xyz.awaitTermination()等待作业完成,或用xyz.stop()停止作业。

内容的提问来源于stack exchange,提问作者zainab

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 18:03:13