Spark 4.0的 MemoryStream是否被移动或更改?
我尝试把我的项目从 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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。