Spark Structured Streaming流流内连接计算滚动平均返回空结果问题
Spark Structured Streaming流流内连接返回空结果排查
我正在用Spark Structured Streaming处理非商用IMDB数据集的流式排名计算,排名公式为(numVotes/averageNumVotes)*averageRating。由于评分是持续流式输入,已将数据集schema修改为带时间戳的事件日志格式。
选择Spark Structured Streaming的原因是熟悉Spark批处理,且官方文档显示框架可自动管理状态无需自定义处理,核心逻辑是先做流式聚合,再执行流流内连接。
遇到的核心问题:流流内连接始终返回空结果,但单独计算两个连接输入时,均存在匹配连接键的行(代码仅水印处理和输出逻辑不同);且相同代码在非流式(批处理)模式下能得到预期的连接结果。
相关演示场景的笔记本(可在VSCode Dev Containers中运行,项目包含安装Spark 3.5.3的devcontainer.json):
- 流内连接返回空结果及调试信息的笔记本
- 无连接的流式输入(包含预期可匹配的连接数据)的笔记本
- 非流式参考实现的笔记本
已尝试将Spark版本从3.5.3升级到3.5.4,问题仍未解决。目前不确定是框架bug还是对官方文档的理解有误,打算在Spark Jira提交bug前寻求社区意见。
内容的提问来源于stack exchange,提问作者alza
相关产品推荐
相关产品推荐

