PyFlink中能否串联使用DataStream API与Table API SQL?
结论
完全支持你描述的开发模式,不需要为每种事件类型单独部署Flink作业。PyFlink 原生打通了 DataStream API 和 Table API/SQL 的转换链路,侧输出分流后的DataStream也可以直接转成Table执行SQL,全流程可以在同一个作业拓扑内完成。
具体实现流程
- 初始化绑定的执行环境:创建
StreamExecutionEnvironment之后,基于该实例初始化对应的StreamTableEnvironment,保证两个环境的拓扑打通。
from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import StreamTableEnvironment, DataTypes from pyflink.datastream.output_tag import OutputTag # 初始化DataStream执行环境 env = StreamExecutionEnvironment.get_execution_environment() # 基于同一个env初始化Table环境,这一步是两个API打通的核心 t_env = StreamTableEnvironment.create(env)
- 接入原始流并完成侧输出分流:通过
addSource接入数据源,完成必要的基础数据清洗转换后,为每种事件类型定义带明确Schema的OutputTag,在处理函数中将不同类型的事件打到对应的侧输出。
# 为不同类型的事件定义侧输出标签,必须明确指定字段名、字段类型 browse_event_tag = OutputTag("browse_event", DataTypes.ROW([ DataTypes.FIELD("user_id", DataTypes.BIGINT()), DataTypes.FIELD("item_id", DataTypes.BIGINT()), DataTypes.FIELD("ts", DataTypes.TIMESTAMP(3)) ])) order_event_tag = OutputTag("order_event", DataTypes.ROW([ DataTypes.FIELD("order_id", DataTypes.STRING()), DataTypes.FIELD("user_id", DataTypes.BIGINT()), DataTypes.FIELD("pay_amount", DataTypes.DECIMAL(10,2)), DataTypes.FIELD("ts", DataTypes.TIMESTAMP(3)) ])) # 接入原始数据源 source_stream = env.add_source(自定义Source实现类) # 编写分流逻辑,在process方法中根据事件类型将数据输出到对应侧输出 processed_stream = source_stream.process(自定义分流处理函数, output_tags=[browse_event_tag, order_event_tag])
- 侧输出流转Table执行SQL:从处理后的流中提取各侧输出的DataStream,直接转换为Table并注册为临时视图,即可针对不同视图编写独立的SQL逻辑。
# 提取各侧输出流 browse_stream = processed_stream.get_side_output(browse_event_tag) order_stream = processed_stream.get_side_output(order_event_tag) # 流转Table并注册为临时视图 t_env.create_temporary_view("browse_event", t_env.from_data_stream(browse_stream)) t_env.create_temporary_view("order_event", t_env.from_data_stream(order_stream)) # 针对不同流执行独立SQL,例如统计分钟级浏览UV、订单GMV browse_uv_stat = t_env.sql_query(""" SELECT TUMBLE_START(ts, INTERVAL '1' MINUTE) as window_start, COUNT(DISTINCT user_id) as uv FROM browse_event GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE) """) order_gmv_stat = t_env.sql_query(""" SELECT TUMBLE_START(ts, INTERVAL '1' MINUTE) as window_start, SUM(pay_amount) as gmv FROM order_event GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE) """) # 为各SQL结果配置对应的Sink输出即可
注意事项
- 定义侧输出
OutputTag时必须明确声明完整的Row类型Schema,否则转Table时会出现类型推断失败、字段无法识别的问题。 - 禁止单独初始化和DataStream环境无绑定关系的TableEnvironment,否则两套API的拓扑无法合并,会提交成独立作业。
- 全流程逻辑编写完成后,仅需调用一次
env.execute()即可提交整个作业,所有分流、转换、SQL计算逻辑会在同一个Flink作业中调度运行。
内容的提问来源于stack exchange,提问作者John
相关产品推荐
相关产品推荐

