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。
步骤如下:
- 创建用于存储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' );
- 将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;
- 对新表执行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
相关产品推荐
相关产品推荐

