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

Flink事件时间时态连接无数据输出问题求助

环境信息

  • Flink版本:1.13
  • Kafka主题分区数:1

左侧探测端(仅追加型DataStream)

数据结构:

{
"eventType": String,
"eventTime": LocalDateTime,
"eventId": String
}

转换为表的代码:

var eventTable = tableEnv.fromDataStream(eventStream, Schema.newBuilder()
.column("eventId", DataTypes.STRING())
.column("eventTime", DataTypes.TIMESTAMP(3))
.column("eventType", DataTypes.STRING())
.watermark("eventTime", $("eventTime"))
.build());

右侧构建端(Kafka+Debezium CDC版本表)

表定义SQL:

CREATE TABLE metadata (
id                  VARCHAR,
eventMetadata       VARCHAR,
origin_ts           TIMESTAMP(3) METADATA FROM 'value.source.timestamp' VIRTUAL,
PRIMARY KEY (id) NOT ENFORCED,
WATERMARK FOR origin_ts AS origin_ts
) WITH (
'connector'                           = 'kafka',
'properties.bootstrap.servers'        = 'SERVER_ADDR',
'properties.group.id'                 = 'SOME_GROUP',
'topic'                               = 'SOME_TOPIC',
'scan.startup.mode'                   = 'latest-offset',
'value.format'                        = 'debezium-json'
)

连接查询语句

SELECT e.eventId, e.eventTime, e.eventType, m.eventMetadata
FROM events_view AS e
JOIN metadata_view FOR SYSTEM_TIME AS OF e.eventTime AS m
ON e.eventId = m.id

已执行的调试操作

  • 设置源空闲超时配置:table.exec.source.idle-timeout -> 5
  • 为水印配置空闲时间,确认水印已正常生成,但数据全部滞留在时态连接表中,始终无输出结果,也未出现任何运行时异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 06:50:49