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

Flink SQL 1.14中Match Recognize消费Join结果报错咨询

问题原因

报错核心是Flink SQL 1.14的MATCH_RECOGNIZE算子不支持消费包含更新(UPDATE)、删除(DELETE)事件的数据流。

你的场景中:

  • EventReference是Postgres CDC表,CDC连接器会输出Upsert流(包含新增、更新、删除事件)
  • Event是Kafka的Append流(仅包含新增事件)
  • 两者执行Inner Join后,输出的数据流会带有Upsert特性(当EventReference有更新/删除操作时,Join结果会产生对应的更新/删除事件),而MATCH_RECOGNIZE无法处理这类流。
解决方案

根据业务需求选择以下方案:

方案1:将EventReference改为静态批量读取(无需实时同步变更)

如果EventReference是静态数据,不需要实时同步Postgres的变更,可以改用JDBC连接器批量读取该表,输出的是Append流,Join后的结果也会是Append流,MATCH_RECOGNIZE即可正常处理。

修改EventReference的创建语句:

CREATE TEMPORARY TABLE EventReference (
      category STRING
     ,commentt STRING
     ,description STRING
     ,type STRING
     ,id STRING PRIMARY KEY
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:postgresql://localhost:5432/postgres',
    'table-name' = 'public.EventReference',
    'username' = 'postgres',
    'password' = 'admin'
);

方案2:将Join结果写入Append流存储后再做模式匹配

如果需要保留EventReference的实时变更同步,可以把Join后的结果写入一个Kafka表(Append模式),再对这个Kafka表执行MATCH_RECOGNIZE。

步骤如下:

  1. 创建用于存储Join结果的Kafka表:
CREATE TABLE Event_Join_Result (
  isFired BOOLEAN,
  eventMrid STRING,
  createDateTime TIMESTAMP(3),
  eventReference_id STRING,
  type STRING,
  WATERMARK FOR createDateTime AS createDateTime - INTERVAL '5' MINUTE
) WITH (
  'connector' = 'kafka',
  'topic' = 'event-join-result',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id' = 'event-join-group',
  'scan.startup.mode' = 'latest-offset',
  'format' = 'json'
);
  1. 将Join结果插入到该表:
INSERT INTO Event_Join_Result
SELECT
  e.isFired,
  e.eventMrid,
  e.createDateTime,
  r.id AS eventReference_id,
  r.type
FROM EventReference r
JOIN Event e ON r.id = e.eventReference_id;
  1. 对新表执行MATCH_RECOGNIZE:
SELECT * from Event_Join_Result
MATCH_RECOGNIZE (
 PARTITION BY eventMrid
 ORDER BY createDateTime
 MEASURES
 FIRST(S.createDateTime) AS firstDateTime
 ONE ROW PER MATCH
 PATTERN (S)
 DEFINE
 S AS S.isFired = False
);

方案3:升级Flink版本

Flink 1.15及以上版本对MATCH_RECOGNIZE的能力做了增强,支持处理Upsert流。如果业务允许,可以考虑升级到更高版本的Flink,直接使用原有SQL逻辑即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 10:45:44