Skip to content

Commit b9bcd7a

Browse files
committed
add incr monitor support for TIS
1 parent 40b55d4 commit b9bcd7a

11 files changed

Lines changed: 274 additions & 122 deletions

File tree

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFactory.java

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,23 +18,33 @@
1818

1919
package com.qlangtech.plugins.incr.flink.cdc.pglike;
2020

21+
import com.alibaba.citrus.turbine.Context;
2122
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2223
import com.qlangtech.plugins.incr.flink.cdc.pglike.PGDTOColValProcess.PGCDCTypeVisitor;
23-
import com.qlangtech.tis.async.message.client.consumer.IConsumerHandle;
2424
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
2525
import com.qlangtech.tis.async.message.client.consumer.impl.MQListenerFactory;
26+
import com.qlangtech.tis.datax.DataXName;
27+
import com.qlangtech.tis.datax.impl.DataxReader;
2628
import com.qlangtech.tis.plugin.annotation.FormField;
2729
import com.qlangtech.tis.plugin.annotation.FormFieldType;
2830
import com.qlangtech.tis.plugin.annotation.Validator;
2931
import com.qlangtech.tis.plugin.ds.DataSourceMeta;
32+
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
33+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
34+
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
35+
import com.qlangtech.tis.util.IPluginContext;
3036

37+
import java.util.List;
3138
import java.util.Objects;
3239

3340
/**
3441
* @author: 百岁(baisui@qlangtech.com)
3542
* @create: 2025-01-19 18:14
3643
**/
3744
public abstract class FlinkCDCPGLikeSourceFactory extends MQListenerFactory {
45+
46+
public static final String FIELD_REPLICA_RULE = "replicaRule";
47+
3848
/**
3949
* The name of the Postgres logical decoding plug-in installed on the server. Supported values are decoderbufs, wal2json, wal2json_rds, wal2json_streaming, wal2json_rds_streaming and pgoutput.
4050
*/
@@ -47,12 +57,12 @@ public abstract class FlinkCDCPGLikeSourceFactory extends MQListenerFactory {
4757
@FormField(ordinal = 1, type = FormFieldType.ENUM, validate = {Validator.require})
4858
public String startupOptions;
4959
// REPLICA IDENTITY
50-
@FormField(ordinal = 2, advance = false, type = FormFieldType.ENUM, validate = {Validator.require})
51-
public String replicaIdentity;
60+
@FormField(ordinal = 2, advance = false, validate = {Validator.require})
61+
public PGLikeReplicaIdentity replicaRule;
5262

5363

54-
public ReplicaIdentity getRepIdentity() {
55-
return ReplicaIdentity.parse(this.replicaIdentity);
64+
public final PGLikeReplicaIdentity getRepIdentity() {
65+
return Objects.requireNonNull(replicaRule, "replicaIdentity can not be null");
5666
}
5767

5868

@@ -72,6 +82,15 @@ public PluginVender getVender() {
7282
return PluginVender.FLINK_CDC;
7383
}
7484

75-
85+
@Override
86+
protected boolean validateMQListenerForm(
87+
IControlMsgHandler msgHandler, Context context, MQListenerFactory sourceFactory) {
88+
FlinkCDCPGLikeSourceFactory incrSource = (FlinkCDCPGLikeSourceFactory) sourceFactory;
89+
DataXName pipeline = msgHandler.getCollectionName();
90+
DataxReader dataxReader = DataxReader.load((IPluginContext) msgHandler, pipeline.getPipelineName());
91+
IDataSourceFactoryGetter dataSourceGetter = (IDataSourceFactoryGetter) dataxReader;
92+
List<ISelectedTab> selectedTabs = dataxReader.getSelectedTabs();
93+
return incrSource.replicaRule.validateSelectedTabs(msgHandler, context, dataSourceGetter, selectedTabs);
94+
}
7695
}
7796
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
/**
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
* <p>
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
* <p>
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
package com.qlangtech.plugins.incr.flink.cdc.pglike;
20+
21+
import com.alibaba.citrus.turbine.Context;
22+
import com.qlangtech.tis.extension.Describable;
23+
import com.qlangtech.tis.extension.TISExtensible;
24+
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
25+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
26+
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
27+
28+
import java.io.Serializable;
29+
import java.util.List;
30+
31+
/**
32+
* PostgreSQL binlog的复制规则,如果需要支持物理删除,必须要使用full,不然更新流程中不会附带需要删除的物理id
33+
*
34+
* @author: 百岁(baisui@qlangtech.com)
35+
* @create: 2024-11-11 16:31
36+
**/
37+
@TISExtensible
38+
public abstract class PGLikeReplicaIdentity implements Describable<PGLikeReplicaIdentity>, Serializable {
39+
protected static final String FULL = "FULL";
40+
protected static final String DEFAULT = "DEFAULT";
41+
42+
public abstract boolean isShallContainBeforeVals();
43+
44+
45+
public abstract boolean validateSelectedTabs(
46+
IControlMsgHandler msgHandler
47+
, Context context
48+
, IDataSourceFactoryGetter dataSourceGetter, List<ISelectedTab> selectedTabs);
49+
}

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/PostgreSQLDeserializationSchema.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@
3939
* @create: 2024-02-21 19:02
4040
**/
4141
public class PostgreSQLDeserializationSchema extends TISDeserializationSchema {
42-
private final ReplicaIdentity replicaIdentity;
42+
private final PGLikeReplicaIdentity replicaIdentity;
4343
private static final Logger logger = LoggerFactory.getLogger(PostgreSQLDeserializationSchema.class);
4444

4545
/**
@@ -50,7 +50,7 @@ public class PostgreSQLDeserializationSchema extends TISDeserializationSchema {
5050
*/
5151
public PostgreSQLDeserializationSchema(List<ISelectedTab> tabs, IFlinkColCreator<FlinkCol> flinkColCreator
5252
, Map<String /*tableName*/, Map<String, Function<RunningContext, Object>>> contextParamValsGetterMapper
53-
, ReplicaIdentity replicaIdentity
53+
, PGLikeReplicaIdentity replicaIdentity
5454
) {
5555
super(new PGDTOColValProcess(tabs, flinkColCreator), new DefaultTableNameConvert(), contextParamValsGetterMapper);
5656
this.replicaIdentity = replicaIdentity;

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/ReplicaIdentity.java

Lines changed: 0 additions & 47 deletions
This file was deleted.
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
package com.qlangtech.plugins.incr.flink.cdc.pglike.replica;
2+
3+
import com.alibaba.citrus.turbine.Context;
4+
import com.qlangtech.plugins.incr.flink.cdc.pglike.PGLikeReplicaIdentity;
5+
import com.qlangtech.tis.extension.Descriptor;
6+
import com.qlangtech.tis.extension.TISExtension;
7+
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
8+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
9+
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
10+
11+
import java.util.List;
12+
13+
/**
14+
*
15+
* @author 百岁 (baisui@qlangtech.com)
16+
* @date 2025/11/23
17+
*/
18+
public class DefaultPGLikeReplicaIdentity extends PGLikeReplicaIdentity {
19+
@Override
20+
public boolean isShallContainBeforeVals() {
21+
return false;
22+
}
23+
24+
@Override
25+
public boolean validateSelectedTabs(IControlMsgHandler msgHandler
26+
, Context context, IDataSourceFactoryGetter dataSourceGetter, List<ISelectedTab> selectedTabs) {
27+
return true;
28+
}
29+
30+
31+
@TISExtension
32+
public static final class DefaultDesc extends Descriptor<PGLikeReplicaIdentity> {
33+
@Override
34+
public String getDisplayName() {
35+
return DEFAULT;
36+
}
37+
}
38+
39+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
package com.qlangtech.plugins.incr.flink.cdc.pglike.replica;
2+
3+
import com.alibaba.citrus.turbine.Context;
4+
import com.qlangtech.plugins.incr.flink.cdc.pglike.PGLikeReplicaIdentity;
5+
import com.qlangtech.tis.extension.Descriptor;
6+
import com.qlangtech.tis.extension.TISExtension;
7+
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
8+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
9+
import com.qlangtech.tis.plugin.ds.JDBCConnection;
10+
import com.qlangtech.tis.plugin.ds.postgresql.PGLikeDataSourceFactory;
11+
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
12+
import org.apache.commons.lang3.StringUtils;
13+
14+
import java.sql.PreparedStatement;
15+
import java.sql.ResultSet;
16+
import java.util.List;
17+
18+
import static com.qlangtech.plugins.incr.flink.cdc.pglike.FlinkCDCPGLikeSourceFactory.FIELD_REPLICA_RULE;
19+
20+
/**
21+
*
22+
* @author 百岁 (baisui@qlangtech.com)
23+
* @date 2025/11/23
24+
*/
25+
public class FullPGLikeReplicaIdentity extends PGLikeReplicaIdentity {
26+
@Override
27+
public boolean isShallContainBeforeVals() {
28+
return true;
29+
}
30+
31+
/**
32+
* FULL: 使用此值需要确保对应的表执行ALTER TABLE your_table_name REPLICA IDENTITY FULL;,表记录更新时会带上更新Before值,使用此方式比较耗费性能。
33+
*
34+
* @param msgHandler
35+
* @param context
36+
* @param dataSourceGetter
37+
* @param selectedTabs
38+
* @return
39+
*/
40+
41+
@Override
42+
public boolean validateSelectedTabs(
43+
IControlMsgHandler msgHandler
44+
, Context context
45+
, IDataSourceFactoryGetter dataSourceGetter, List<ISelectedTab> selectedTabs) {
46+
47+
PGLikeDataSourceFactory ds = (PGLikeDataSourceFactory) dataSourceGetter.getDataSourceFactory();
48+
49+
boolean[] validate = new boolean[]{true};
50+
ds.visitFirstConnection((conn) -> {
51+
validate[0] = validateReplicaIdentity(msgHandler, context, conn, ds.tabSchema, selectedTabs);
52+
});
53+
54+
return validate[0];
55+
}
56+
57+
/**
58+
* 验证所有选中的表是否配置了REPLICA IDENTITY FULL
59+
*
60+
* @param msgHandler 消息处理器
61+
* @param context 上下文
62+
* @param conn 数据库连接
63+
* @param selectedTabs 选中的表列表
64+
* @return 验证是否通过
65+
*/
66+
private boolean validateReplicaIdentity(IControlMsgHandler msgHandler, Context context,
67+
JDBCConnection conn, final String tabSchema, List<ISelectedTab> selectedTabs) {
68+
if (StringUtils.isEmpty(tabSchema)) {
69+
throw new IllegalStateException("param tableSchema can not be empty");
70+
}
71+
boolean allValid = true;
72+
73+
// SQL查询:检查表的REPLICA IDENTITY设置
74+
// relreplident: 'd' = default, 'f' = full, 'i' = index, 'n' = nothing
75+
String sql = "SELECT c.relreplident " +
76+
"FROM pg_class c " +
77+
"JOIN pg_namespace n ON c.relnamespace = n.oid " +
78+
"WHERE n.nspname = ? AND c.relname = ?";
79+
// String schema;
80+
try (PreparedStatement stmt = conn.preparedStatement(sql)) {
81+
for (ISelectedTab tab : selectedTabs) {
82+
String tableName = tab.getName();
83+
84+
// 查询表的REPLICA IDENTITY设置
85+
stmt.setString(1, tabSchema);
86+
stmt.setString(2, tableName);
87+
88+
try (ResultSet rs = stmt.executeQuery()) {
89+
if (rs.next()) {
90+
String replIdent = rs.getString("relreplident");
91+
92+
// 检查是否为FULL模式
93+
if (!"f".equals(replIdent)) {
94+
String fullTableName = tabSchema + "." + tableName;
95+
String errMessage = "表'" + fullTableName + "'未配置REPLICA IDENTITY FULL," +
96+
"请执行:ALTER TABLE " + fullTableName + " REPLICA IDENTITY FULL;";
97+
msgHandler.addFieldError(context, FIELD_REPLICA_RULE, errMessage);
98+
allValid = false;
99+
return allValid;
100+
}
101+
} else {
102+
// 表不存在
103+
throw new IllegalStateException("table'" + tabSchema + "." + tableName + "' is not exist");
104+
}
105+
}
106+
}
107+
} catch (Exception e) {
108+
throw new RuntimeException("the validator of replica is vailed", e);
109+
}
110+
111+
return allValid;
112+
}
113+
114+
115+
@TISExtension
116+
public static final class DefaultDesc extends Descriptor<PGLikeReplicaIdentity> {
117+
@Override
118+
public String getDisplayName() {
119+
return FULL;
120+
}
121+
122+
123+
}
124+
}

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/resources/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFactory.json

Lines changed: 3 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,7 @@
11
{
2-
"replicaIdentity": {
3-
"label": "REPLICA IDENTITY",
4-
"dftVal": "DEFAULT",
5-
"enum": [
6-
{
7-
"val": "DEFAULT",
8-
"label": "DEFAULT"
9-
},
10-
{
11-
"val": "FULL",
12-
"label": "FULL"
13-
}
14-
]
2+
"replicaRule": {
3+
"label": "Replica Rule",
4+
"dftVal": "FULL"
155
},
166
"decodingPluginName": {
177
"dftVal": "decoderbufs",

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/resources/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFactory.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,15 @@ Debezium startup options
1010
* `Latest`:
1111
Never to perform snapshot on the monitored database tables upon first startup, just read from the end of the binlog which means only have the changes since the connector was started.
1212

13-
## replicaIdentity
13+
## replicaRule
1414

1515
在 PostgreSQL 中,ALTER TABLE ... REPLICA IDENTITY 命令用于指定在逻辑复制或行级触发器中如何标识已更新或删除的行。https://developer.aliyun.com/ask/575334
1616

1717
可选项有以下两个
1818
* `FULL`: 使用此值需要确保对应的表执行`ALTER TABLE your_table_name REPLICA IDENTITY FULL;`,表记录更新时会带上更新Before值,使用此方式比较耗费性能。
1919
* `DEFAULT`: 默认值,更新删除操作时不会带上Before值。
2020

21+
如目标端需要实现`物理删除`,必须选择`FULL`
22+
2123

2224

0 commit comments

Comments
 (0)