Skip to content

Commit 1c572dd

Browse files
committed
add TISPBReporter for metric reporter
1 parent 320bfbb commit 1c572dd

10 files changed

Lines changed: 224 additions & 9 deletions

File tree

tis-datax/tis-datax-common-plugin/src/main/java/com/qlangtech/tis/plugin/datax/common/RdbmsReaderContext.java

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

2121
import com.qlangtech.tis.datax.IDataxReaderContext;
22-
import com.qlangtech.tis.plugin.datax.common.RdbmsReaderContext.ISplitTableContext;
2322
import com.qlangtech.tis.plugin.ds.DataSourceFactory;
2423
import com.qlangtech.tis.plugin.ds.IDataSourceDumper;
2524
import org.apache.commons.lang.StringUtils;

tis-datax/tis-datax-common-plugin/src/main/java/com/qlangtech/tis/plugin/ds/BasicDataSourceFactory.java

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -62,7 +62,7 @@
6262
* @create: 2021-06-06 19:48
6363
**/
6464
public abstract class BasicDataSourceFactory extends DataSourceFactory
65-
implements JdbcUrlBuilder, IPluginStore.AfterPluginSaved, Describable.IRefreshable, IDBAuthorizeTokenGetter {
65+
implements JdbcUrlBuilder, IPluginStore.AfterPluginSaved, Describable.IRefreshable, IDBAuthorizeTokenGetter, SplitTableStrategyAbility {
6666
private static String TYPE_NAME_JSON = "json";
6767

6868
public static Optional<String> MYSQL_ESCAPE_COL_CHAR = Optional.of("`");
@@ -93,8 +93,6 @@ public static boolean isJSONColumnType(DataType type) {
9393
public String password;
9494

9595

96-
97-
9896
/**
9997
* 数据库编码
10098
*/
@@ -127,6 +125,10 @@ public String getPassword() {
127125
return this.password;
128126
}
129127

128+
@Override
129+
public SplitTableStrategy getSplitTableStrategy() {
130+
throw new UnsupportedOperationException();
131+
}
130132

131133
@Override
132134
public List<ColumnMetaData> getTableMetadata(boolean inSink, IPluginContext pluginContext, final EntityName table) {

tis-datax/tis-datax-doris-plugin/src/main/java/com/qlangtech/tis/plugin/ds/doris/DorisSourceFactory.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,8 @@
3939
import com.qlangtech.tis.plugin.ds.DataType.DefaultTypeVisitor;
4040
import com.qlangtech.tis.plugin.ds.JDBCConnection;
4141
import com.qlangtech.tis.plugin.ds.JDBCTypes;
42+
import com.qlangtech.tis.plugin.ds.NoneSplitTableStrategy;
43+
import com.qlangtech.tis.plugin.ds.SplitTableStrategy;
4244
import com.qlangtech.tis.plugin.ds.TableNotFoundException;
4345
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
4446
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
@@ -88,6 +90,15 @@ public class DorisSourceFactory extends BasicDataSourceFactory {
8890
@FormField(ordinal = 8, type = FormFieldType.TEXTAREA, validate = {Validator.require})
8991
public String loadUrl;
9092

93+
@Override
94+
public SplitTableStrategy getSplitTableStrategy() {
95+
NoneSplitTableStrategy splitTableStrategy = new NoneSplitTableStrategy();
96+
if (StringUtils.isEmpty(this.nodeDesc)) {
97+
throw new IllegalStateException("nodeDesc can not be null");
98+
}
99+
splitTableStrategy.host = this.nodeDesc;
100+
return splitTableStrategy;
101+
}
91102

92103
public List<String> getLoadUrls() {
93104
return DorisSourceFactory.getLoadUrls(this.loadUrl);

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -221,8 +221,7 @@ private static Thread monitorDistributeCommand(Integer jobId, DataXJobInfo jobIn
221221
public void run() {
222222
UpdateCounterMap status = new UpdateCounterMap();
223223
status.setFrom(NetUtils.getHost());
224-
logger.info("start to listen the dataX job taskId:{},jobName:{},dataXName:{} overseer cancel", jobId,
225-
jobInfo, dataXName);
224+
logger.info("start to listen the dataX job taskId:{},jobName:{},dataXName:{} overseer cancel", jobId, jobInfo, dataXName);
226225
TableSingleDataIndexStatus dataXStatus = new TableSingleDataIndexStatus();
227226
dataXStatus.setUUID(jobInfo.jobFileName);
228227
status.addTableCounter(IAppSourcePipelineController.DATAX_FULL_PIPELINE + dataXName, dataXStatus);
Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
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.tis.plugin.doris;
20+
21+
import com.qlangtech.tis.extension.TISExtension;
22+
import com.qlangtech.tis.extension.impl.IOUtils;
23+
import com.qlangtech.tis.plugin.annotation.FormField;
24+
import com.qlangtech.tis.plugin.annotation.FormFieldType;
25+
import com.qlangtech.tis.plugin.annotation.Validator;
26+
import com.qlangtech.tis.plugin.datax.RdbmsDataxContext;
27+
import com.qlangtech.tis.plugin.datax.SelectedTab;
28+
import com.qlangtech.tis.plugin.datax.common.BasicDataXRdbmsReader;
29+
import com.qlangtech.tis.plugin.datax.common.RdbmsReaderContext;
30+
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
31+
import com.qlangtech.tis.plugin.ds.IDataSourceDumper;
32+
33+
/**
34+
* @author: 百岁(baisui@qlangtech.com)
35+
* @create: 2025-07-04 05:50
36+
**/
37+
public class DataXDorisReader extends BasicDataXRdbmsReader<BasicDataSourceFactory> {
38+
public static final String DATAX_NAME = "Doris";
39+
40+
@FormField(ordinal = 1, type = FormFieldType.ENUM, validate = {Validator.require, Validator.identity})
41+
public Boolean splitPk;
42+
43+
44+
public static String getDftTemplate() {
45+
return IOUtils.loadResourceFromClasspath(DataXDorisReader.class, "mysql-reader-tpl.vm");
46+
}
47+
48+
@Override
49+
protected RdbmsReaderContext createDataXReaderContext(
50+
String jobName, SelectedTab tab, IDataSourceDumper dumper) {
51+
BasicDataSourceFactory dsFactory = this.getDataSourceFactory();
52+
53+
RdbmsDataxContext rdbms = new RdbmsDataxContext(this.dataXName);
54+
rdbms.setJdbcUrl(dumper.getDbHost());
55+
rdbms.setUsername(dsFactory.getUserName());
56+
rdbms.setPassword(dsFactory.getPassword());
57+
return new DorisDataXReaderContext(jobName, tab.getName(), rdbms, this);
58+
}
59+
60+
@TISExtension()
61+
public static class DefaultDescriptor extends BasicDataXRdbmsReaderDescriptor {
62+
public DefaultDescriptor() {
63+
super();
64+
}
65+
66+
@Override
67+
public boolean isSupportIncr() {
68+
return false;
69+
}
70+
71+
@Override
72+
public String getDisplayName() {
73+
return DATAX_NAME;
74+
}
75+
76+
@Override
77+
public EndType getEndType() {
78+
return EndType.Doris;
79+
}
80+
}
81+
}
Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,92 @@
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.tis.plugin.doris;
20+
21+
import com.qlangtech.tis.plugin.datax.RdbmsDataxContext;
22+
import com.qlangtech.tis.plugin.datax.common.RdbmsReaderContext;
23+
import com.qlangtech.tis.plugin.datax.common.RdbmsReaderContext.ISplitTableContext;
24+
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
25+
import com.qlangtech.tis.plugin.ds.SplitTableStrategy;
26+
import org.apache.commons.lang.StringUtils;
27+
28+
import java.util.List;
29+
import java.util.Objects;
30+
31+
/**
32+
* @author: 百岁(baisui@qlangtech.com)
33+
* @create: 2025-07-04 05:52
34+
**/
35+
public class DorisDataXReaderContext
36+
extends RdbmsReaderContext<DataXDorisReader, BasicDataSourceFactory> implements ISplitTableContext {
37+
private final RdbmsDataxContext rdbmsContext;
38+
private final SplitTableStrategy splitTableStrategy;
39+
40+
public DorisDataXReaderContext(String name, String sourceTableName
41+
, RdbmsDataxContext mysqlContext, DataXDorisReader dataXReader) {
42+
super(name, sourceTableName, null, Objects.requireNonNull(dataXReader, "dataXReader can not be null"));
43+
this.rdbmsContext = mysqlContext;
44+
this.splitTableStrategy = Objects.requireNonNull(dsFactory.getSplitTableStrategy(), "splitTableStrategy can not be null");
45+
}
46+
47+
/**
48+
* 是否执行分表导入
49+
*
50+
* @return
51+
*/
52+
@Override
53+
public boolean isSplitTable() {
54+
// this.splitTableStrategy.getAllPhysicsTabs(this.dsFactory, this.getJdbcUrl(), this.sourceTableName);
55+
// return !(this.splitTableStrategy instanceof NoneSplitTableStrategy);
56+
return this.splitTableStrategy.isSplittable();
57+
}
58+
59+
/**
60+
* 分表列表
61+
*
62+
* @return
63+
*/
64+
@Override
65+
public String getSplitTabs() {
66+
List<String> allPhysicsTabs = this.splitTableStrategy.getAllPhysicsTabs(dsFactory, this.getJdbcUrl(), this.sourceTableName);
67+
return getEntitiesWithQuotation(allPhysicsTabs);
68+
}
69+
70+
public String getDataXName() {
71+
return rdbmsContext.getDataXName();
72+
}
73+
74+
public String getTabName() {
75+
return rdbmsContext.getTabName();
76+
}
77+
78+
public String getPassword() {
79+
return rdbmsContext.getPassword();
80+
}
81+
82+
public String getUsername() {
83+
return rdbmsContext.getUsername();
84+
}
85+
86+
public String getJdbcUrl() {
87+
if (StringUtils.isEmpty(rdbmsContext.getJdbcUrl())) {
88+
throw new NullPointerException("rdbmsContext.getJdbcUrl() can not be empty");
89+
}
90+
return rdbmsContext.getJdbcUrl();
91+
}
92+
}

tis-datax/tis-ds-mysql-plugin/src/main/java/com/qlangtech/tis/plugin/ds/mysql/MySQLDataSourceFactory.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,10 +105,13 @@ protected TableInDB createTableInDB() {
105105
+ "relevant prop splitTableStrategy can not be null").createTableInDB(this);
106106
}
107107

108+
@Override
109+
public SplitTableStrategy getSplitTableStrategy() {
110+
return Objects.requireNonNull(splitTableStrategy);
111+
}
108112

109113
@Override
110114
public List<String> getAllPhysicsTabs(DataXJobSubmit.TableDataXEntity tabEntity) {
111-
// return super.getAllPhysicsTabs(tabEntity);
112115
return this.splitTableStrategy.getAllPhysicsTabs(this, tabEntity);
113116
}
114117

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
{
2+
"dbName": {
3+
"enum": "com.qlangtech.tis.util.PluginItems.getExistDbs(\"Doris\")",
4+
"creator": {
5+
"plugin": [
6+
{
7+
"descName": "Doris"
8+
}
9+
]
10+
}
11+
},
12+
"template": {
13+
"dftVal": "com.qlangtech.tis.plugin.datax.DataxMySQLReader.getDftTemplate()"
14+
}
15+
}

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

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,18 +18,25 @@
1818

1919
package com.qlangtech.plugins.incr.flink.metrics;
2020

21+
import org.apache.flink.metrics.CharacterFilter;
22+
import org.apache.flink.metrics.Counter;
2123
import org.apache.flink.metrics.MetricConfig;
24+
import org.apache.flink.metrics.MetricGroup;
2225
import org.apache.flink.metrics.reporter.AbstractReporter;
2326
import org.apache.flink.metrics.reporter.Scheduled;
2427

28+
import java.util.Map;
29+
2530
/**
31+
* https://github.com/datavane/tis/issues/397
32+
*
2633
* @author: 百岁(baisui@qlangtech.com)
2734
* @create: 2025-05-14 10:24
2835
**/
2936
public class TISPBReporter extends AbstractReporter implements Scheduled {
3037
@Override
3138
public String filterCharacters(String s) {
32-
return "";
39+
return CharacterFilter.NO_OP_FILTER.filterCharacters(s);
3340
}
3441

3542
@Override
@@ -45,5 +52,11 @@ public void close() {
4552
@Override
4653
public void report() {
4754

55+
for (Map.Entry<Counter, String> entry : counters.entrySet()) {
56+
Counter counter = entry.getKey();
57+
String metricName = entry.getValue();
58+
System.out.println(metricName + ": " + counter.getCount());
59+
}
60+
4861
}
4962
}

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,6 @@
3232
public class TISPBReporterFactory implements MetricReporterFactory {
3333
@Override
3434
public MetricReporter createMetricReporter(Properties properties) {
35-
return null;
35+
return new TISPBReporter();
3636
}
3737
}

0 commit comments

Comments
 (0)