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

Spark Structured Streaming:MemoryStream新增数据无法输出至控制台问题

问题描述

原本需测试Spark结构化流从流式数据源读取并写入S3/Parquet的pipeline,简化后实现将MemoryStream新增记录输出至控制台的功能。当前代码存在问题:首次添加的"Alice"和"Bob"可正常输出至控制台,但在超过查询触发间隔的sleep操作后,新增的"Mouse"无法触发第二批次输出。

实际输出:

Batch: 0

+-----+
|value|
+-----+
|Alice|
| Bob|
+-----+

期望看到包含"Mouse"的第二批次输出,但运行代码后从未出现,相关代码如下:

import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
import org.apache.spark.sql.{Encoder, Encoders, SparkSession}
import org.apache.spark.sql.execution.streaming.MemoryStream
import org.apache.spark.sql.streaming.{OutputMode, Trigger}
import org.apache.spark.sql.types.{IntegerType, StringType, StructType}

object StreamTableDemo {
  def main(args: Array[String]): Unit = {
    implicit val spark = SparkSession
      .builder
      .appName("StreamTableDemo")
      .master("local[*]")
      .getOrCreate()

    // Define the schema for the stream data
    val schema = new StructType() .add("id", IntegerType)

    implicit val stringEncoder = ExpressionEncoder[String]

    // Create a MemoryStream and a DataFrame based on the schema
    val stream = new MemoryStream[String](2, spark.sqlContext, Some(1))
    val streamDF = stream.toDF()


    // Create a temporary view from the stream DataFrame
    streamDF.createOrReplaceTempView("T")

    // Define a streaming query to output the records in T to the console
    val query = spark.table("T")
      .writeStream
      .outputMode(OutputMode.Append())
      .trigger(Trigger.ProcessingTime("1 seconds"))
      .format("console")
      .start()

    // Add data to the stream
    stream.addData("Alice")
    stream.addData("Bob")

    // Wait for the query to terminate
    query.awaitTermination()
    Thread.sleep(2 * 1000)
    stream.addData("Mouse")
  }
}
问题原因与解决方法
  • 核心问题:阻塞方法执行顺序错误
    query.awaitTermination()是阻塞式方法,会一直等待流式查询终止才会执行后续代码。当前代码中,这条语句放在了Thread.sleep和stream.addData("Mouse")之前,导致后续添加"Mouse"的代码根本没有机会执行。

  • 解决步骤

    1. 调整代码执行顺序,将query.awaitTermination()移至所有数据添加操作之后,确保"Mouse"能被添加到流中;
    2. 移除未被使用的schema定义(当前MemoryStream[String]并未用到该schema,属于冗余代码)。

修改后的代码如下:

import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
import org.apache.spark.sql.{SparkSession}
import org.apache.spark.sql.execution.streaming.MemoryStream
import org.apache.spark.sql.streaming.{OutputMode, Trigger}

object StreamTableDemo {
  def main(args: Array[String]): Unit = {
    implicit val spark = SparkSession
      .builder
      .appName("StreamTableDemo")
      .master("local[*]")
      .getOrCreate()

    implicit val stringEncoder = ExpressionEncoder[String]

    // Create a MemoryStream and a DataFrame
    val stream = new MemoryStream[String](2, spark.sqlContext, Some(1))
    val streamDF = stream.toDF()

    // Create a temporary view from the stream DataFrame
    streamDF.createOrReplaceTempView("T")

    // Define a streaming query to output the records in T to the console
    val query = spark.table("T")
      .writeStream
      .outputMode(OutputMode.Append())
      .trigger(Trigger.ProcessingTime("1 seconds"))
      .format("console")
      .start()

    // Add initial data
    stream.addData("Alice")
    stream.addData("Bob")

    // Wait for 2 seconds then add new data
    Thread.sleep(2 * 1000)
    stream.addData("Mouse")

    // Wait for the query to terminate
    query.awaitTermination()
  }
}
  • 说明
    修改后,主线程先添加初始数据,等待2秒后添加"Mouse",最后才进入阻塞等待查询终止。由于触发间隔设置为1秒,添加"Mouse"后,流式查询会在下一个触发周期处理该数据,控制台就会输出包含"Mouse"的第二批次结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 05:55:34