【问题标题】:Spark Structured Streaming Unit Test in JavaJava中的Spark结构化流单元测试
【发布时间】:2021-02-23 02:57:03
【问题描述】:

我正在开发一个 api,用于从 Kafka 读取数据并使用 Spark 结构化流将数据写入 Java 中的 blob 存储。我找不到为此编写单元测试的方法。我有一个返回数据集的读取器类和一个将数据集作为输入并以指定格式写入 blob 存储的写入器类。我在 MemoryStream 上看到了一些博客,但我认为这还不够。

提前致谢。

【问题讨论】:

  • 您需要更精确一点:您要对代码的哪一部分进行单元测试?从 kafka 读取,写入 blob 存储(哪种类型?),或者您可能正在应用的一些中间转换?根据部分的不同,单元测试可能会有所不同。对于输入(kafka)和输出(存储),您需要模拟外部系统进行测试

标签: java apache-spark junit apache-kafka spark-structured-streaming


【解决方案1】:

显然,您可以参考这个答案,了解我们如何使用内存流进行单元测试 - Unit Test - structured streaming

此外,您还可以查看 Holden Karau 的这个火花测试基地。 Spark testing base

您可以模拟来自 Kafka 的流式数据帧,并针对该数据帧之上的代码中的转换运行测试用例。

示例:

static Dataset<Row> createTestStreamingDataFrame() {
    MemoryStream<String> testStream= new MemoryStream<String>(100, sqlContext(), Encoders.STRING());
    testStream.addData((Arrays.asList("1,1","2,2","3,3")).toSeq());
    return testStream.toDF().selectExpr(
        "cast(split(value,'[,]')[0] as int) as testCol1",
        "cast(split(value,'[,]')[1] as int) as testCol2");
}

【讨论】:

  • @mhlaskar1991 - 这是否澄清了您的疑问?!
猜你喜欢
  • 2017-05-04
  • 1970-01-01
  • 2018-01-18
  • 1970-01-01
  • 2017-03-06
  • 1970-01-01
  • 2020-01-31
  • 1970-01-01
  • 2020-11-30
相关资源
最近更新 更多