Skip to content

Commit 64d4a57

Browse files
committed
add ontology for TIS datavane/tis#495
1 parent 5d5c8ba commit 64d4a57

14 files changed

Lines changed: 217 additions & 103 deletions

File tree

pom.xml

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,7 @@
104104
<!-- <hudi.version>0.12.2</hudi.version>-->
105105
<!-- <hive.version>3.1.3</hive.version>-->
106106

107-
<hudi.version>0.14.1</hudi.version>
107+
108108

109109
<hive.version>3.1.3</hive.version>
110110

@@ -924,11 +924,7 @@
924924
<artifactId>access-modifier-checker</artifactId>
925925
<version>${access-modifier-checker.version}</version>
926926
</plugin>
927-
<!-- <plugin>-->
928-
<!-- <groupId>com.mycila</groupId>-->
929-
<!-- <artifactId>license-maven-plugin</artifactId>-->
930-
<!-- <version>4.1</version>-->
931-
<!-- </plugin>-->
927+
<!---->
932928
<plugin> <!-- not gated by Incrementals profiles, since we want the incrementalify goal to be available from the start -->
933929
<groupId>io.jenkins.tools.incrementals</groupId>
934930
<artifactId>incrementals-maven-plugin</artifactId>

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -161,7 +161,7 @@ static void addTransformerInfo(Set<TransformerInfo> tinfos,
161161
pluginCtx = IPluginContext.namedContext(new DataXName(pipeName, resType));
162162
}
163163
Key transformerRuleKey = TransformerRuleKey.createStoreKey(pluginCtx, resType, pipeName, "dump");
164-
XmlFile sotre = transformerRuleKey.getSotreFile();
164+
XmlFile sotre = transformerRuleKey.getStoreXmlFile();
165165
File parent = sotre.getFile().getParentFile();
166166
if (!parent.exists()) {
167167
// return Collections.emptySet();

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

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -78,28 +78,28 @@ public static boolean isJSONColumnType(DataType type) {
7878
*/
7979
@FormField(ordinal = 1, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.hostWithoutPort})
8080
public String nodeDesc;
81-
@FormField(ordinal = 2, type = FormFieldType.INT_NUMBER, validate = {Validator.require, Validator.integer})
81+
@FormField(ordinal = 4, type = FormFieldType.INT_NUMBER, validate = {Validator.require, Validator.integer})
8282
public int port;
8383
// 数据库名称
84-
@FormField(ordinal = 3, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.identity})
84+
@FormField(ordinal = 5, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.identity})
8585
public String dbName;
8686

87-
@FormField(ordinal = 5, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.user_name})
87+
@FormField(ordinal = 7, prompt4llm = true, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.user_name})
8888
public String userName;
8989

90-
@FormField(ordinal = 7, type = FormFieldType.PASSWORD, validate = {Validator.none_blank, Validator.require})
90+
@FormField(ordinal = 9, prompt4llm = true, type = FormFieldType.PASSWORD, validate = {Validator.none_blank, Validator.require})
9191
public String password;
9292

9393

9494
/**
9595
* 数据库编码
9696
*/
97-
@FormField(ordinal = 13, type = FormFieldType.ENUM, validate = {Validator.require, Validator.identity})
97+
@FormField(ordinal = 15, type = FormFieldType.ENUM, validate = {Validator.require, Validator.identity})
9898
public String encode;
9999
/**
100100
* 附加参数
101101
*/
102-
@FormField(ordinal = 15, advance = true, type = FormFieldType.INPUTTEXT)
102+
@FormField(ordinal = 17, advance = true, type = FormFieldType.INPUTTEXT)
103103
public String extraParams;
104104

105105
public static <DS extends DataSourceFactory> DS getDs(String dbName) {
@@ -125,7 +125,7 @@ public String getPassword() {
125125

126126
@Override
127127
public SplitTableStrategy getSplitTableStrategy() {
128-
throw new UnsupportedOperationException();
128+
throw new UnsupportedOperationException();
129129
}
130130

131131
@Override

tis-datax/tis-datax-common-plugin/src/test/java/com/qlangtech/tis/plugin/datax/TestDefaultDataxProcessor.java

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -19,12 +19,11 @@
1919
package com.qlangtech.tis.plugin.datax;
2020

2121
import com.google.common.collect.Lists;
22+
import com.qlangtech.tis.datax.StoreResourceType;
2223
import com.qlangtech.tis.datax.TableAlias;
23-
import com.qlangtech.tis.datax.TableAliasMapper;
2424
import com.qlangtech.tis.extension.impl.XmlFile;
2525
import com.qlangtech.tis.manage.IAppSource;
2626
import com.qlangtech.tis.plugin.KeyedPluginStore;
27-
import com.qlangtech.tis.datax.StoreResourceType;
2827
import com.qlangtech.tis.plugin.common.PluginDesc;
2928
import com.qlangtech.tis.plugin.test.BasicTest;
3029
import org.apache.commons.io.FileUtils;
@@ -68,23 +67,23 @@ public void testSaveProcess() {
6867
assertEquals(dataxProcessor.dptId, loadDataxProcessor.dptId);
6968
assertEquals(dataxProcessor.recept, loadDataxProcessor.recept);
7069

71-
TableAliasMapper tabAlias1 = loadDataxProcessor.getTabAlias(null);
72-
assertEquals(1, tabAlias1.size());
73-
74-
tabAlias1.forEach((key, val) -> {
75-
assertEquals(tabAlias.getFrom(), key);
76-
77-
assertEquals(tabAlias.getFrom(), val.getFrom());
78-
assertEquals(tabAlias.getTo(), val.getTo());
79-
});
70+
// TableAliasMapper tabAlias1 = loadDataxProcessor.getTabAlias(null);
71+
// assertEquals(1, tabAlias1.size());
72+
//
73+
// tabAlias1.forEach((key, val) -> {
74+
// assertEquals(tabAlias.getFrom(), key);
75+
//
76+
// assertEquals(tabAlias.getFrom(), val.getFrom());
77+
// assertEquals(tabAlias.getTo(), val.getTo());
78+
// });
8079

8180
// for (Map.Entry<String, TableAlias> entry : tabAlias1.entrySet()) {
8281
//
8382
// }
8483
} finally {
8584
try {
8685
KeyedPluginStore.AppKey appKey = new KeyedPluginStore.AppKey(null, StoreResourceType.parse(false), appName, IAppSource.class);
87-
XmlFile storeFile = appKey.getSotreFile();
86+
XmlFile storeFile = appKey.getStoreXmlFile();
8887
FileUtils.forceDelete(storeFile.getFile().getParentFile());
8988
} catch (IOException e) {
9089
throw new RuntimeException(e);
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
package com.qlangtech.tis.plugin.ds;
2+
3+
import com.qlangtech.tis.extension.Describable;
4+
5+
/**
6+
*
7+
* @author 百岁 (baisui@qlangtech.com)
8+
* @date 2026/4/13
9+
*/
10+
public abstract class DataSourceCatalog implements Describable<DataSourceCatalog> {
11+
public abstract void appendJdbcUrl(StringBuffer jdbcUrl, String dbName);
12+
// {
13+
// jdbcUrl.append(dbName);
14+
// }
15+
16+
public abstract String getFullTableName(String dbName, String tableName);
17+
18+
19+
}

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

Lines changed: 42 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,14 @@
3333
import com.qlangtech.tis.plugin.annotation.Validator;
3434
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
3535
import com.qlangtech.tis.plugin.ds.DBConfig;
36+
import com.qlangtech.tis.plugin.ds.DataSourceCatalog;
3637
import com.qlangtech.tis.plugin.ds.DataType;
3738
import com.qlangtech.tis.plugin.ds.DataType.DefaultTypeVisitor;
3839
import com.qlangtech.tis.plugin.ds.JDBCConnection;
3940
import com.qlangtech.tis.plugin.ds.JDBCTypes;
4041
import com.qlangtech.tis.plugin.ds.NoneSplitTableStrategy;
4142
import com.qlangtech.tis.plugin.ds.SplitTableStrategy;
43+
import com.qlangtech.tis.plugin.ds.TableInDB;
4244
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
4345
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
4446
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
@@ -84,9 +86,17 @@ public class DorisSourceFactory extends BasicDataSourceFactory {
8486
}
8587
}
8688

89+
@FormField(ordinal = 4, validate = {Validator.require})
90+
public DataSourceCatalog catalog;
91+
8792
@FormField(ordinal = 8, type = FormFieldType.TEXTAREA, validate = {Validator.require})
8893
public String loadUrl;
8994

95+
@Override
96+
protected TableInDB createTableInDB() {
97+
return TableInDB.create(this, (tab) -> catalog.getFullTableName(this.dbName, tab));
98+
}
99+
90100
@Override
91101
public SplitTableStrategy getSplitTableStrategy() {
92102
NoneSplitTableStrategy splitTableStrategy = new NoneSplitTableStrategy();
@@ -112,7 +122,7 @@ public String buidJdbcUrl(DBConfig db, String ip, String dbName) {
112122
jdbcUrl.append("jdbc:mysql://").append(ip).append(":").append(this.port);
113123

114124
if (StringUtils.isNotEmpty(dbName)) {
115-
jdbcUrl.append("/").append(dbName);
125+
catalog.appendJdbcUrl(jdbcUrl.append("/"), dbName);
116126
}
117127
return jdbcUrl.toString();
118128
}
@@ -254,38 +264,38 @@ protected boolean validateDSFactory(final IControlMsgHandler msgHandler, final C
254264
try {
255265
Boolean success = HttpUtils.get(new URL(clusterInfoApiUrl.toString()) //
256266
, new PostFormStreamProcess<Boolean>( //
257-
ConfigFileContext.setAuthorizationHeader( dorisDS.getUserName(), dorisDS.getPassword())) { //
258-
@Override
259-
public ContentType getContentType() {
260-
return ContentType.JSON;
261-
}
262-
263-
@Override
264-
public Duration getSocketReadTimeout() {
265-
return socketReadTimeout;
266-
}
267-
268-
@Override
269-
public Boolean p(int status, InputStream stream, Map<String, List<String>> headerFields) {
270-
try {
271-
JSONObject result = JSONObject.parseObject(IOUtils.toString(stream, TisUTF8.get()));
272-
final String msg = result.getString("msg");
273-
if (!"success".equals(msg)) {
274-
msgHandler.addFieldError(context, FIELD_KEY_LOAD_URL, msg);
275-
return false;
267+
ConfigFileContext.setAuthorizationHeader(dorisDS.getUserName(), dorisDS.getPassword())) { //
268+
@Override
269+
public ContentType getContentType() {
270+
return ContentType.JSON;
271+
}
272+
273+
@Override
274+
public Duration getSocketReadTimeout() {
275+
return socketReadTimeout;
276+
}
277+
278+
@Override
279+
public Boolean p(int status, InputStream stream, Map<String, List<String>> headerFields) {
280+
try {
281+
JSONObject result = JSONObject.parseObject(IOUtils.toString(stream, TisUTF8.get()));
282+
final String msg = result.getString("msg");
283+
if (!"success".equals(msg)) {
284+
msgHandler.addFieldError(context, FIELD_KEY_LOAD_URL, msg);
285+
return false;
286+
}
287+
} catch (IOException e) {
288+
throw new RuntimeException(e);
289+
}
290+
return true;
291+
}
292+
293+
@Override
294+
public void error(int status, InputStream errstream, IOException e) throws Exception {
295+
logger.warn(e.getMessage(), e);
296+
msgHandler.addFieldError(context, FIELD_KEY_LOAD_URL, IOUtils.toString(errstream, TisUTF8.get()));
276297
}
277-
} catch (IOException e) {
278-
throw new RuntimeException(e);
279-
}
280-
return true;
281-
}
282-
283-
@Override
284-
public void error(int status, InputStream errstream, IOException e) throws Exception {
285-
logger.warn(e.getMessage(), e);
286-
msgHandler.addFieldError(context, FIELD_KEY_LOAD_URL, IOUtils.toString(errstream, TisUTF8.get()));
287-
}
288-
});
298+
});
289299
if (success == null || !success) {
290300
break;
291301
}
Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
package com.qlangtech.tis.plugin.ds.impl;
2+
3+
import com.qlangtech.tis.extension.Descriptor;
4+
import com.qlangtech.tis.extension.TISExtension;
5+
import com.qlangtech.tis.plugin.ds.DataSourceCatalog;
6+
7+
/**
8+
* 默认的catalog ,在jdbc连接时候不使用特定的catalog值
9+
*
10+
* @author 百岁 (baisui@qlangtech.com)
11+
* @date 2026/4/13
12+
*/
13+
public class CatalogDefault extends DataSourceCatalog {
14+
@Override
15+
public void appendJdbcUrl(StringBuffer jdbcUrl, String dbName) {
16+
// super.appendJdbcUrl(jdbcUrl, dbName);
17+
jdbcUrl.append(dbName);
18+
}
19+
20+
@Override
21+
public String getFullTableName(String dbName, String tableName) {
22+
return tableName;
23+
}
24+
25+
@TISExtension
26+
public static class DefaultDesc extends Descriptor<DataSourceCatalog> {
27+
public DefaultDesc() {
28+
super();
29+
}
30+
31+
@Override
32+
public String getDisplayName() {
33+
return SWITCH_DEFAULT;
34+
}
35+
}
36+
}
Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
1+
package com.qlangtech.tis.plugin.ds.impl;
2+
3+
import com.qlangtech.tis.extension.Descriptor;
4+
import com.qlangtech.tis.extension.TISExtension;
5+
import com.qlangtech.tis.plugin.annotation.FormField;
6+
import com.qlangtech.tis.plugin.annotation.FormFieldType;
7+
import com.qlangtech.tis.plugin.annotation.Validator;
8+
import com.qlangtech.tis.plugin.ds.DataSourceCatalog;
9+
10+
/**
11+
*
12+
* @author 百岁 (baisui@qlangtech.com)
13+
* @date 2026/4/13
14+
*/
15+
public class CatalogSpecific extends DataSourceCatalog {
16+
@FormField(ordinal = 1, type = FormFieldType.INPUTTEXT, validate = {Validator.require, Validator.db_col_name})
17+
public String name;
18+
19+
@Override
20+
public void appendJdbcUrl(StringBuffer jdbcUrl, String dbName) {
21+
jdbcUrl.append(name).append(".").append(dbName);
22+
}
23+
24+
@Override
25+
public String getFullTableName(String dbName, String tableName) {
26+
return this.name + "." + dbName + "." + tableName;
27+
}
28+
29+
@TISExtension
30+
public static class DefaultDesc extends Descriptor<DataSourceCatalog> {
31+
public DefaultDesc() {
32+
super();
33+
}
34+
35+
@Override
36+
public String getDisplayName() {
37+
return SWITCH_CUSTOMIZE;
38+
}
39+
}
40+
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,9 @@
2020
"port": {
2121
"dftVal": 9030
2222
},
23+
"catalog": {
24+
"dftVal": "default"
25+
},
2326
"loadUrl": {
2427
"help": "Doris FE的地址用于Streamload,可以为多个fe地址,fe_ip:fe_http_port",
2528
"dftVal": "[]",

tis-datax/tis-datax-doris-plugin/src/test/java/com/qlangtech/tis/plugin/datax/doris/TestDataXDorisWriter.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
import com.qlangtech.tis.plugin.ds.doris.DorisSourceFactory;
4242
import com.qlangtech.tis.plugin.ds.doris.TestDorisSourceFactory;
4343
import com.qlangtech.tis.trigger.util.JsonUtil;
44+
import com.qlangtech.tis.util.DefaultDescriptorsJSON;
4445
import com.qlangtech.tis.util.DescriptorsJSON;
4546
import junit.framework.TestCase;
4647
import org.apache.commons.io.FileUtils;
@@ -140,7 +141,7 @@ public void testDescriptorsJSONGenerate() {
140141
EasyMock.replay(dataxReader);
141142
DataXDorisWriter writer = new DataXDorisWriter();
142143

143-
DescriptorsJSON descJson = new DescriptorsJSON(writer.getDescriptor());
144+
DescriptorsJSON descJson = new DefaultDescriptorsJSON(writer.getDescriptor());
144145

145146
JsonUtil.assertJSONEqual(DataXDorisWriter.class, "doris-datax-writer-descriptor.json"
146147
, descJson.getDescriptorsJSON(), (m, e, a) -> {

0 commit comments

Comments
 (0)