【发布时间】:2019-10-14 08:25:00
【问题描述】:
我有两个流并希望在窗口内将第二个流加入第一个流,因为我需要对与会话相关的两个流的连接进行一些计算(流中的一个控制会话)。
实际上,从文档中可以看出,(会话)窗口只允许在单个流上进行计算,而不是在连接中。
我尝试使用会话窗口和协处理器功能的组合,但结果与我的预期不完全一致。
有没有办法在 Flink 中合并与会话窗口相关的两个流?
【问题讨论】:
标签: session join stream apache-flink
我有两个流并希望在窗口内将第二个流加入第一个流,因为我需要对与会话相关的两个流的连接进行一些计算(流中的一个控制会话)。
实际上,从文档中可以看出,(会话)窗口只允许在单个流上进行计算,而不是在连接中。
我尝试使用会话窗口和协处理器功能的组合,但结果与我的预期不完全一致。
有没有办法在 Flink 中合并与会话窗口相关的两个流?
【问题讨论】:
标签: session join stream apache-flink
Flink 的 DataStream API 包含一个会话窗口连接,描述为 here。
你必须看看它的语义是否符合你的想法。会话间隙由在该时间间隔内没有事件的两个流定义,并且连接是内部连接,因此如果会话窗口仅包含来自一个流的元素,则不会发出任何输出。
如果这不能满足您的需求,那么我建议使用 CoProcessFunction,但没有会话窗口。换句话说,我建议您自己实现所有逻辑。
【讨论】: