Spark 4.0的 MemoryStream是否被移动或更改?

编程语言 2026-07-11

我尝试把我的项目从 Spark 3.5 升级到 Spark 4.0。在此过程中,我在单元测试中遇到了这个问题。

error: cannot find symbol
  import org.apache.spark.sql.execution.streaming.MemoryStream;

error: cannot find symbol
          MemoryStream inputStream =

以下是导致问题的代码片段

SQLContext spark = TestSparkUtils.getSparkSession().sqlContext();

MemoryStream inputStream =
        new MemoryStream<Row>(
                1,
                TestSparkUtils.getSparkSession().sqlContext(),
                Option.apply(1),
                ExpressionEncoder.apply(TRANSFORM_SCHEMA));

Dataset<Row> df =
        inputStream
                .toDS()
                .toDF(
                        "kafkaKey",
                        "t_uid",
                        "timestamp",
                        "resource_name",
                        "metrics",
                        "attributes");
Seq<Row> collection = JavaConverters.asScalaIteratorConverter(rows.iterator()).asScala().toSeq();

inputStream.addData(collection);

这个类在 Spark 4.0 中是否已经被弃用?如果是,是否有我可以使用的替代类?抱歉,这个问题可能有点傻,但我对Spark还算新手。

解决方案

spark v3.5.8 MemoryStream 已在文件 scala/org/apache/spark/sql/execution/streaming/memory.scala 中定义。

对于 spark 4.0.0,看起来该文件(以及其他一些文件)已从 org.apache.spark.sql.execution.streaming 移动到 org.apache.spark.sql.execution.streaming.runtime,你可以在 Commit c57556c SPARK-52787 Reorganize streaming execution dir around runtime and checkpoint areas 看到

sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/memory.scala -> sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/runtime/memory.scala
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章