Apache Flink - 端到端测试如何终止输入源

Xel*_*eli 4 integration-testing end-to-end data-stream apache-flink flink-streaming

我在批处理中使用 apache flink 一段时间了,但现在我们想将此批处理作业转换为流作业。我遇到的问题是如何运行端到端测试。

它如何在批处理作业中工作

使用批处理时,我们使用 Cucumber 创建端到端测试。

  • 我们将填充我们读取的 hbase 表
  • 运行批处理作业
  • 等待它完成
  • 验证结果

流作业中的问题

我们希望对流作业执行类似的操作,但流作业并未真正完成。

所以:

  • 填充我们读取的消息队列
  • 运行流作业。
  • 等待它完成(如何?)
  • 验证结果

我们可以在每次测试后等待 5 秒,并假设所有内容都已处理完毕,但这会大大减慢一切速度。

问题:

有哪些方法或最佳实践可以在流式 Flink 作业上运行端到端测试,而不会在 x 秒后强制终止 Flink 作业

Dav*_*son 5

大多数 Flink DataStream 源,如果从有限输入读取,将在到达末尾时注入值为 LONG.MAX_VALUE 的水印,之后作业将被终止。

Flink训练练习展示了一种对 Flink 作业进行端到端测试的方法。我建议克隆github 存储库并查看测试的设置方式。他们使用自定义源和接收器并重定向输入和输出以进行测试。

文档中也对此主题进行了一些讨论。