如何用PyFlink将三个含关联ID的数据流合并为一个?
PyFlink多流关联与状态更新实现方案疑问
我在使用PyFlink的DataStream API做数据转换时遇到了问题,也可以接受使用Table API的方案。当前PyFlink作业通过Kafka连接器从三个不同主题读取数据流,已完成JSON消息到自定义Python数据类的解析、字段过滤和嵌套列表扁平化操作。
数据流结构
# C流(有限流,每条消息c_id唯一) C = {'c_id': 'c_1', … } # M流(有限流,每条消息m_id唯一,关联C的c_id) M = {'m_id': 'm_1', 'c_id': 'c_1', … } # O流(持续流,每条消息o_id唯一,关联M的m_id) O = {'o_id': 'o_1', 'm_id': 'm_1', … }
期望输出结构
Output = {'c_id': 'c_1', 'm_ids': ['m_1', 'm_4', … ], 'o_ids': ['o_1', 'o_6', … ] }
业务约束
- C是有限数据流,每条消息包含唯一
c_id; - M是更大的有限数据流,每条消息包含唯一
m_id,且关联C中的某个c_id; - O是持续数据流,每条消息包含唯一
o_id,且关联M中的某个m_id; - 需要为C流的每条消息生成一个输出对象,存储匹配的M和O消息,并随新O消息持续更新,最终输出到Kafka。
非流处理的参考实现
(注:原代码中if m.m_id == c.c_id:应为笔误,修正为if m.c_id == c.c_id以符合业务逻辑)
# Join C and M output = [] for c in C: out = {'c_id': c.c_id, 'm_ids': [], 'o_ids': []} for m in M: if m.c_id == c.c_id: out['m_ids'].append(m.m_id) output.append(out) # Join O for o in O: for out in output: if o.m_id in out['m_ids']: out['o_ids'].append(o.o_id)
已尝试的方案
- 使用Broadcast State将C流作为MapState传入M流,通过
KeyedBroadcastProcessFunction匹配添加m_id; - 使用
CoFlatMapFunction连接两个流,但该方案需要先处理完第一个流才能处理第二个,不符合流处理的实时性需求。
疑问
- 是否适合使用窗口连接实现该需求?
- PyFlink的Python API是否支持对应的窗口连接?
- 该场景下的最优实现方案是什么?
内容的提问来源于stack exchange,提问作者Matthew Weaver
相关产品推荐
相关产品推荐

