PyFlink TableAPI多源归一到中间处理表print连接器报错如何解决
问题根因
你遇到的报错是因为print连接器仅支持作为数据输出sink使用,不支持作为可读取的数据源。你给中间处理表process指定了print连接器,后续需要从该表读取数据做计算时自然触发校验失败。
你的中间表仅作为多源合并后的临时计算载体,不需要关联任何外部存储,完全不需要定义带connector的物理表,用临时视图即可实现需求。
解决方案
- 删掉原代码中定义
process表的那段带print连接器的CREATE TABLE语句 - 对两个数据源转换后的结果做
UNION ALL合并,再注册为临时视图,后续就可以直接用process名称做后续计算,用法和你之前的逻辑完全兼容
修正后关键代码
首先修正UDF注册逻辑,你原代码存在重复注册的bug:
source_udf = udf(sourceUdf(), result_type=DataTypes.ROW([DataTypes.FIELD('recordKey', DataTypes.STRING()),DataTypes.FIELD('tm', DataTypes.TIMESTAMP(3)), DataTypes.FIELD('content', DataTypes.STRING()) ])) t_env.register_function("sourceUdf", source_udf) two_udf = udf(source_two_Udf(), result_type=DataTypes.ROW([DataTypes.FIELD('recordKey', DataTypes.STRING()),DataTypes.FIELD('tm', DataTypes.TIMESTAMP(3)), DataTypes.FIELD('content', DataTypes.STRING()) ])) # 原代码这里重复注册为sourceUdf,修正为单独的名字 t_env.register_function("sourceUdf2", two_udf) sink_udf = udf(sinkUdf(), result_type=DataTypes.STRING()) t_env.register_function("sinkUdf", sink_udf)
然后处理两个源并注册中间临时视图:
# 处理第一个源,注意原路径多了空格要删掉 tbl1 = t_env.from_path('source_table') tbl1 = tbl1.map(source_udf).alias('recordKey', 'tm', 'content') # 处理第二个源 tbl2 = t_env.from_path('second_source_table') tbl2 = tbl2.map(two_udf).alias('recordKey', 'tm', 'content') # 合并两个源的数据 process_tbl = tbl1.union_all(tbl2) # 注册为临时视图,后续可以直接用process名称访问 t_env.create_temporary_view('process', process_tbl)
后续计算逻辑不需要改动,和你原有写法完全兼容:
merge = t_env.from_path('process') merge = merge.map(sink_udf).alias('entry') merge = merge.insert_into('sink_table') t_env.execute("test")
其他注意事项
如果需要在中间表上定义水位线做窗口计算,可以在合并后的process_tbl上直接定义:
process_tbl = process_tbl.select( col('recordKey'), col('tm'), col('content') ).watermark_by('tm', col('tm') - lit(5).second)
再注册为临时视图即可。
内容的提问来源于stack exchange,提问作者Paul
相关产品推荐
相关产品推荐

