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"的代码根本没有机会执行。解决步骤
- 调整代码执行顺序,将
query.awaitTermination()移至所有数据添加操作之后,确保"Mouse"能被添加到流中; - 移除未被使用的
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
相关产品推荐
相关产品推荐

