Skip to content

Commit 44bccba

Browse files
committed
add moidfy KingBaseDataSourceFactory
1 parent 75fd613 commit 44bccba

20 files changed

Lines changed: 297 additions & 123 deletions

File tree

tis-datax/tis-datax-kingbase-plugin/src/main/java/com/qlangtech/tis/plugin/ds/kingbase/KingBaseDataSourceFactory.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
package com.qlangtech.tis.plugin.ds.kingbase;
2020

2121
import com.alibaba.citrus.turbine.Context;
22+
import com.kingbase8.KBProperty;
2223
import com.qlangtech.tis.extension.TISExtension;
2324
import com.qlangtech.tis.lang.TisException;
2425
import com.qlangtech.tis.plugin.annotation.FormField;
@@ -70,6 +71,10 @@ protected java.sql.Driver createDriver() {
7071
@Override
7172
protected java.util.Properties extractSetJdbcProps(java.util.Properties props) {
7273
Objects.requireNonNull(this.dispatch, "dispatch can not be null").extractSetJdbcProps(props);
74+
if (StringUtils.isEmpty(this.encode)) {
75+
throw new IllegalStateException("param encode can not be empty");
76+
}
77+
props.setProperty(KBProperty.CLIENT_ENCODING.getName(), this.encode);
7378
return props;
7479
}
7580

@@ -178,6 +183,7 @@ protected boolean validateAll(IControlMsgHandler msgHandler, Context context, Po
178183

179184
// return super.validateAll(msgHandler, context, postFormVals);
180185
}
186+
181187
@Override
182188
public Optional<String> getDefaultDataXReaderDescName() {
183189
return Optional.of(KingBaseDataSourceFactory.KingBase_NAME);

tis-datax/tis-datax-local-executor-utils/src/main/java/com/qlangtech/tis/datax/DataxExecutor.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@
5353
import com.qlangtech.tis.realtime.transfer.TableSingleDataIndexStatus;
5454
import com.qlangtech.tis.realtime.utils.NetUtils;
5555
import com.qlangtech.tis.realtime.yarn.rpc.MasterJob;
56+
import com.qlangtech.tis.realtime.yarn.rpc.PipelineFlinkTaskId;
5657
import com.qlangtech.tis.realtime.yarn.rpc.UpdateCounterMap;
5758
import com.tis.hadoop.rpc.RpcServiceReference;
5859
import com.tis.hadoop.rpc.StatusRpcClientFactory;
@@ -223,7 +224,7 @@ public void run() {
223224
logger.info("start to listen the dataX job taskId:{},jobName:{},dataXName:{} overseer cancel", jobId, jobInfo, dataXName);
224225
TableSingleDataIndexStatus dataXStatus = new TableSingleDataIndexStatus();
225226
dataXStatus.setUUID(jobInfo.jobFileName);
226-
status.addTableCounter(IAppSourcePipelineController.DATAX_FULL_PIPELINE + dataXName, dataXStatus);
227+
status.setPipelineTableCounterMetric(new PipelineFlinkTaskId(dataXName, IAppSourcePipelineController.DATAX_FULL_PIPELINE), dataXStatus);
227228

228229
while (true) {
229230
status.setUpdateTime(System.currentTimeMillis());

tis-datax/tis-datax-mongodb-plugin/src/main/java/com/qlangtech/tis/plugin/datax/mongo/MongoColValGetter.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
package com.qlangtech.tis.plugin.datax.mongo;
2020

2121
import org.bson.BsonDocument;
22+
import org.bson.BsonType;
2223
import org.bson.BsonValue;
2324

2425
import java.time.ZoneId;
@@ -60,7 +61,7 @@ public FlinkPropValGetter(MongoCMeta cmeta, ZoneId zone) {
6061
@Override
6162
public Object apply(BsonDocument document) {
6263
BsonValue val = cmeta.getBsonVal(document);
63-
if (val == null) {
64+
if (val == null || (val.getBsonType() == BsonType.NULL)) {
6465
return null;
6566
}
6667
return MongoDataXColUtils.createCol(cmeta, val, false, zone);

tis-incr/tis-flink-cdc-kingbase-plugin/src/test/java/com/qlangtech/plugins/incr/flink/cdc/kingbase/source/KingBaseConnectionPoolFactoryTest.java

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -58,16 +58,11 @@ public void testCreatePooledDataSource() throws Exception {
5858

5959
}
6060

61-
// @Test
62-
// public void testGetSchemaList() {
63-
// //
64-
//
65-
// PostgresSourceConfig cfg = new StubPostgresSourceConfig();
66-
// List<String> schemaList = KingBaseConnectionPoolFactory.getSchemaList(cfg);
67-
// Assert.assertTrue(CollectionUtils.isEqualCollection(Lists.newArrayList("schema1", "schema2"), schemaList));
68-
// }
69-
//
70-
//
61+
62+
/**
63+
* int subtaskId, StartupOptions startupOptions, List<String> databaseList, List<String> schemaList, List<String> tableList, int splitSize, int splitMetaGroupSize, double distributionFactorUpper, double distributionFactorLower, boolean includeSchemaChanges, boolean closeIdleReaders, Properties dbzProperties, Configuration dbzConfiguration, String driverClassName, String hostname, int port, String username, String password, int fetchSize, String serverTimeZone, Duration connectTimeout, int connectMaxRetries, int connectionPoolSize, @Nullable String chunkKeyColumn, boolean skipSnapshotBackfill, boolean isScanNewlyAddedTableEnabled, int lsnCommitCheckpointsDelay, boolean assignUnboundedChunkFirst
64+
*/
65+
7166
private static class StubPostgresSourceConfig extends PostgresSourceConfig {
7267
public StubPostgresSourceConfig() {
7368
super(0, null, Lists.newArrayList("test"), Lists.newArrayList("public"), null, 0, 0, 0d, 0d, true, true, null, null, null, "192.168.28.201", 4321, "kingbase", "123456", 0, "Asia/Shanghai", Duration.ofSeconds(10), 1, 1, null, true, true);

tis-incr/tis-flink-cdc-oracle-plugin/pom.xml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,10 @@
7171
<groupId>org.apache.flink</groupId>
7272
<artifactId>flink-connector-debezium</artifactId>
7373
</exclusion>
74+
<exclusion>
75+
<groupId>io.debezium</groupId>
76+
<artifactId>debezium-core</artifactId>
77+
</exclusion>
7478
</exclusions>
7579
</dependency>
7680

tis-incr/tis-flink-cdc-oracle-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/oracle/FlinkCDCOracleSourceFactory.java

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,8 @@
2020

2121
import com.google.common.collect.Lists;
2222
import com.qlangtech.plugins.incr.debuzium.DebuziumPropAssist;
23-
2423
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2524
import com.qlangtech.tis.annotation.Public;
26-
import com.qlangtech.tis.async.message.client.consumer.IConsumerHandle;
2725
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
2826
import com.qlangtech.tis.async.message.client.consumer.IMQListener;
2927
import com.qlangtech.tis.async.message.client.consumer.impl.MQListenerFactory;
@@ -36,14 +34,10 @@
3634
import com.qlangtech.tis.plugin.ds.DataSourceMeta;
3735
import io.debezium.config.Field;
3836
import io.debezium.connector.oracle.OracleConnectorConfig;
39-
import io.debezium.connector.oracle.OracleConnectorConfig.LogMiningStrategy;
4037
import org.apache.commons.lang3.tuple.Triple;
4138
import org.apache.flink.cdc.connectors.base.options.StartupOptions;
42-
import org.apache.kafka.common.config.ConfigDef.Importance;
43-
import org.apache.kafka.common.config.ConfigDef.Width;
4439

4540
import java.util.List;
46-
import java.util.Objects;
4741
import java.util.function.BooleanSupplier;
4842
import java.util.function.Function;
4943

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/cdc/BiFunction.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,4 +27,14 @@
2727
**/
2828
public abstract class BiFunction implements Function<Object, Object>, Serializable, DeFunction {
2929

30+
/**
31+
* 在Transformer中使用getString方法取值
32+
* //@see AbstractTransformerRecord.getString(String field, boolean origin)
33+
* @param val
34+
* @return
35+
*/
36+
public String toStringVal(Object val) {
37+
38+
return val != null ? String.valueOf(val) : null;
39+
}
3040
}

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/cdc/FlinkCDCPipelineEventProcess.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,10 @@ public org.apache.flink.cdc.common.types.DataType getDataType() {
4646
return this.dataType;
4747
}
4848

49+
@Override
50+
public String toStringVal(Object val) {
51+
return fieldProcessDelegate.toStringVal(val);
52+
}
4953

5054
public static class FlinkPipelineStringConvert extends BiFunction {
5155
@Override
@@ -63,7 +67,6 @@ public Object apply(Object o) {
6367
}
6468

6569

66-
6770
public static class FlinkPipelineDecimalConvert extends BiFunction {
6871
// private final DataType type;
6972

@@ -97,5 +100,10 @@ public Object apply(Object o) {
97100
LocalDateTime v = (LocalDateTime) super.apply(o);
98101
return org.apache.flink.cdc.common.data.TimestampData.fromLocalDateTime(v);
99102
}
103+
104+
@Override
105+
public String toStringVal(Object val) {
106+
return datetimeFormatter.format(((org.apache.flink.cdc.common.data.TimestampData) val).toLocalDateTime());
107+
}
100108
}
101109
}

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/cdc/FlinkCol.java

Lines changed: 31 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@
2323
import com.qlangtech.tis.plugin.ds.ISelectedTab;
2424
import com.qlangtech.tis.realtime.SelectedTableTransformerRules;
2525
import com.qlangtech.tis.util.IPluginContext;
26-
import org.apache.commons.beanutils.converters.DateTimeConverter;
2726
import org.apache.commons.lang3.StringUtils;
2827
import org.apache.flink.table.data.RowData;
2928
import org.apache.flink.table.types.DataType;
@@ -41,6 +40,7 @@
4140
import java.util.function.Function;
4241
import java.util.stream.Collectors;
4342

43+
4444
/**
4545
* @author: 百岁(baisui@qlangtech.com)
4646
* @create: 2022-02-17 16:11
@@ -90,14 +90,30 @@ public static <T extends IColMetaGetter> List<FlinkCol> getAllTabColsMeta(List<T
9090
public static List<FlinkCol> createSourceCols(IPluginContext pluginContext
9191
, final ISelectedTab tab, IFlinkColCreator<FlinkCol> sourceFlinkColCreator
9292
, Optional<SelectedTableTransformerRules> transformerOpt) {
93+
return createCols(pluginContext, tab, sourceFlinkColCreator, transformerOpt, (rules) -> {
94+
return rules.originColsWithContextParamsFlinkCol();
95+
});
96+
}
97+
98+
/**
99+
*
100+
* @param pluginContext
101+
* @param tab
102+
* @param sourceFlinkColCreator
103+
* @param transformerOpt
104+
* @param transformerColOverwrite 使用flink-cdc pipeline模式同步模式下,由于直接使用source表的col meta 来映射 sink端表的列类型,需要使用 rules.overwriteColsWithContextParams()
105+
* @return
106+
*/
107+
public static List<FlinkCol> createCols(IPluginContext pluginContext
108+
, final ISelectedTab tab, IFlinkColCreator<FlinkCol> sourceFlinkColCreator
109+
, Optional<SelectedTableTransformerRules> transformerOpt, Function<SelectedTableTransformerRules, List<FlinkCol>> transformerColOverwrite) {
93110
List<FlinkCol> sourceColsMeta = null;
94111
if (transformerOpt.isPresent()) {
95112
SelectedTableTransformerRules rules = transformerOpt.get();
96-
sourceColsMeta = rules.originColsWithContextParamsFlinkCol();
113+
sourceColsMeta = transformerColOverwrite.apply(rules);
97114
} else {
98115
sourceColsMeta = getAllTabColsMeta(tab.getCols(), sourceFlinkColCreator);
99116
}
100-
101117
return sourceColsMeta;
102118
}
103119

@@ -170,27 +186,6 @@ public FlinkCol setPk(boolean pk) {
170186
return this;
171187
}
172188

173-
public Object processVal(DTOConvertTo convertTo, Object val) {
174-
if (val == null) {
175-
return null;
176-
}
177-
return convertTo.targetValGetter.apply(this, val);
178-
}
179-
180-
public enum DTOConvertTo {
181-
RowData((flinkCol, val) -> {
182-
return flinkCol.rowDataProcess.apply(val);
183-
}),
184-
FlinkCDCPipelineEvent((flinkCol, val) -> {
185-
return flinkCol.flinkCDCPipelineEventProcess.apply(val);
186-
});
187-
188-
private final java.util.function.BiFunction<FlinkCol, Object, Object> targetValGetter;
189-
190-
private DTOConvertTo(java.util.function.BiFunction<FlinkCol, Object, Object> targetValGetter) {
191-
this.targetValGetter = targetValGetter;
192-
}
193-
}
194189

195190
public static BiFunction ByteBuffer() {
196191
return new ByteBufferProcess();
@@ -265,6 +260,7 @@ public Object apply(Object o) {
265260

266261
public static class LocalDateProcess extends BiFunction {
267262
public final static DateTimeFormatter dateFormatter = DateTimeFormatter.ofPattern("yyyy-M-d");
263+
public final static DateTimeFormatter dateFormatterFull = DateTimeFormatter.ofPattern("yyyy-MM-dd");
268264

269265
@Override
270266
public Object apply(Object o) {
@@ -275,21 +271,31 @@ public Object apply(Object o) {
275271
return (LocalDate) o;
276272
}
277273

274+
@Override
275+
public String toStringVal(Object val) {
276+
return dateFormatterFull.format((LocalDate) val);
277+
}
278+
278279
@Override
279280
public Object deApply(Object o) {
280281
return dateFormatter.format((LocalDate) o);
281282
}
282283
}
283284

284285
public static class DateTimeProcess extends BiFunction {
285-
private final static DateTimeFormatter datetimeFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
286+
protected final static DateTimeFormatter datetimeFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
286287
private final static DateTimeFormatter datetimeFormatter_with_zone = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss'Z'");
287288

288289
@Override
289290
public Object deApply(Object o) {
290291
return datetimeFormatter.format((LocalDateTime) o);
291292
}
292293

294+
@Override
295+
public String toStringVal(Object val) {
296+
return datetimeFormatter.format((LocalDateTime) val);
297+
}
298+
293299
@Override
294300
public Object apply(Object o) {
295301
if (o instanceof String) {

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/metrics/TISPBReporter.java

Lines changed: 14 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import com.qlangtech.tis.realtime.transfer.IIncreaseCounter;
2626
import com.qlangtech.tis.realtime.transfer.TableSingleDataIndexStatus;
2727
import com.qlangtech.tis.realtime.yarn.rpc.MasterJob;
28+
import com.qlangtech.tis.realtime.yarn.rpc.PipelineFlinkTaskId;
2829
import com.qlangtech.tis.realtime.yarn.rpc.UpdateCounterMap;
2930
import com.tis.hadoop.rpc.RpcServiceReference;
3031
import com.tis.hadoop.rpc.StatusRpcClientFactory;
@@ -164,27 +165,27 @@ public void report() {
164165
for (Map.Entry<Counter, String> entry : counters.entrySet()) {
165166
Counter counter = entry.getKey();
166167
String metricID = entry.getValue();
167-
System.out.println(metricID + ": " + counter.getCount());
168+
// System.out.println(metricID + ": " + counter.getCount());
168169
metricGroup = metricIdentifierMapper.get(metricID);
169170
if (metricGroup != null) {
170-
System.out.println(metricGroup);
171+
// System.out.println(metricGroup);
171172

172-
metrics.add(new UseableMetricForTIS(entry.getKey(), metricGroup.getKey(), metricGroup.getRight()));
173+
metrics.add(new UseableMetricForTIS(counter, /**metricName*/metricGroup.getKey(), metricGroup.getRight()));
173174
}
174175
}
175176

176-
sendMetric2TISAssemble(metrics);
177+
this.sendMetric2TISAssemble(metrics);
177178
}
178179

179180
private void sendMetric2TISAssemble(List<UseableMetricForTIS> metrics) {
180181
// 汇总一个索引中所有focus table的增量信息
181-
Map<String, TableSingleDataIndexStatus> tabCounterMapper = Maps.newHashMap();
182+
Map<PipelineFlinkTaskId, TableSingleDataIndexStatus> tabCounterMapper = Maps.newHashMap();
182183

183-
String pipelineName = null;
184+
PipelineFlinkTaskId pipelineName = null;
184185
TableSingleDataIndexStatus singleDataIndexStatus = null;
185186
String host = null;
186187
for (UseableMetricForTIS metric : metrics) {
187-
pipelineName = metric.getPipelineName();
188+
pipelineName = new PipelineFlinkTaskId(metric.getPipelineName(), metric.getFlinkTaskId());
188189
if (host == null) {
189190
host = metric.getHost();
190191
}
@@ -195,26 +196,22 @@ private void sendMetric2TISAssemble(List<UseableMetricForTIS> metrics) {
195196
tabCounterMapper.put(pipelineName, singleDataIndexStatus);
196197
}
197198
singleDataIndexStatus.put(metric.metricName, metric.counter.getCount());
198-
// IncrCounter tableIncrCounter = new
199-
// IncrCounter((int)entry.getValue().getIncreasePastLast());
200-
// tableIncrCounter.setAccumulationCount(entry.getValue().getAccumulation());
201-
// tableUpdateCounter.put(entry.getKey(), tableIncrCounter);
202-
// 只记录一个消费总量和当前时间
203199
}
204200

205201
if (MapUtils.isNotEmpty(tabCounterMapper)) {
206-
UpdateCounterMap updateCounterMap = new UpdateCounterMap();
202+
UpdateCounterMap pipelineUpdateCounterMap = new UpdateCounterMap();
207203
if (StringUtils.isEmpty(host)) {
208204
throw new IllegalStateException("host can not be empty");
209205
}
210-
updateCounterMap.setFrom(host);
206+
pipelineUpdateCounterMap.setFrom(host);
211207
tabCounterMapper.forEach((tisPipeline, tabCounter) -> {
212-
updateCounterMap.addTableCounter(tisPipeline, tabCounter);
208+
209+
pipelineUpdateCounterMap.setPipelineTableCounterMetric(tisPipeline, tabCounter);
213210
});
214211
/**
215212
* 服务端:IncrStatusUmbilicalProtocolImpl
216213
*/
217-
getRpcService().reportStatus(updateCounterMap);
214+
getRpcService().reportStatus(pipelineUpdateCounterMap);
218215
}
219216
}
220217

@@ -232,7 +229,7 @@ String getHost() {
232229
}
233230

234231
String getFlinkTaskId() {
235-
return metricGroup.getAllVariables().get(ScopeFormat.SCOPE_TASK_VERTEX_ID);
232+
return metricGroup.getAllVariables().get(ScopeFormat.SCOPE_TASK_SUBTASK_INDEX);
236233
}
237234

238235
public UseableMetricForTIS(Counter counter, String metricName, MetricGroup metricGroup) {

0 commit comments

Comments
 (0)