You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Flink Table API消费两个源流并实现联合窗口聚合

核心实现逻辑

你需要的和DataStream union完全一致的语义,直接用Table API的UNION ALL语法即可实现,整个流程分为4步:

  1. 对齐两个事件流的Schema:两个流的字段数量、字段顺序、字段类型必须完全一致,同时事件时间字段、水位线生成策略也要保持统一,避免合并后水位线推进异常。
  2. 分别消费两个源流,注册为Table对象(可以通过DDL注册源表,也可以从DataStream转成Table)。
  3. 执行UNION ALL合并两个Table。
  4. 在合并后的Table上执行窗口聚合操作,逻辑和单流窗口聚合完全一致。

代码示例

from pyflink.table import EnvironmentSettings, TableEnvironment
from pyflink.table.expressions import col, lit
from pyflink.table.window import Tumble

# 1. 初始化流处理表环境
env_settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_settings)

# 2. 注册两个数据源表,Schema完全对齐,均定义事件时间和水位线
t_env.execute_sql("""
    CREATE TABLE source_1 (
        user_id STRING,
        pay_amount BIGINT,
        event_time TIMESTAMP(3),
        WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'topic_1',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'json'
    )
""")

t_env.execute_sql("""
    CREATE TABLE source_2 (
        user_id STRING,
        pay_amount BIGINT,
        event_time TIMESTAMP(3),
        WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'topic_2',
        'properties.bootstrap.servers' = 'localhost:9092',
        'format' = 'json'
    )
""")

# 3. 用UNION ALL合并两个流
union_table = t_env.sql_query("""
    SELECT user_id, pay_amount, event_time FROM source_1
    UNION ALL
    SELECT user_id, pay_amount, event_time FROM source_2
""")

# 4. 对合并后的流执行滚动窗口聚合
result_table = union_table.window(
    Tumble.over(lit(10).minutes).on(col("event_time")).alias("win")
).group_by(
    col("win"), col("user_id")
).select(
    col("user_id"),
    col("win").start.alias("window_start"),
    col("win").end.alias("window_end"),
    col("pay_amount").sum.alias("total_pay")
)

# 输出结果(示例为打印到控制台,可替换为实际业务的sink逻辑)
result_table.execute().print()

注意事项

  • 如果两个原始流的Schema不一致,需要先通过select调整字段别名、顺序、类型后再执行UNION ALL,比如第二个流的金额字段叫amount,就需要调整为SELECT user_id, amount as pay_amount, event_time FROM source_2
  • 禁止使用UNION替代UNION ALL,UNION会触发全局去重逻辑,会产生额外的状态开销,且和DataStream union的全量合并语义不符
  • 两个源的水位线容忍时间建议配置为相同值,避免合并后水位线被较慢的源卡住

内容的提问来源于stack exchange,提问作者Alfred

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 03:54:06