Skip to content

Commit 2587652

Browse files
committed
add wildcard pattern support for kafka tableName match ,issue:datavane/tis#468
1 parent cd8811e commit 2587652

40 files changed

Lines changed: 307 additions & 978 deletions

File tree

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

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@
4343
import com.qlangtech.tis.extension.impl.SuFormProperties;
4444
import com.qlangtech.tis.extension.impl.XmlFile;
4545
import com.qlangtech.tis.manage.IAppSource;
46+
import com.qlangtech.tis.plugin.IEndTypeGetter;
4647
import com.qlangtech.tis.plugin.IPluginStore;
4748
import com.qlangtech.tis.plugin.IPluginStore.AfterPluginSaved;
4849
import com.qlangtech.tis.plugin.IdentityName;
@@ -82,7 +83,6 @@
8283
import java.util.Optional;
8384
import java.util.Set;
8485
import java.util.function.BiFunction;
85-
import java.util.function.Function;
8686
import java.util.stream.Collectors;
8787

8888
/**
@@ -421,13 +421,18 @@ public DataXCfgGenerator.GenerateCfgs getDataxCfgFileNames(IPluginContext plugin
421421

422422

423423
@TISExtension()
424-
public static class DescriptorImpl extends Descriptor<IAppSource> {
424+
public static class DescriptorImpl extends Descriptor<IAppSource> implements IEndTypeGetter {
425425

426426
public DescriptorImpl() {
427427
super();
428428
this.registerSelectOptions(DefaultDataxProcessor.KEY_FIELD_NAME, () -> ParamsConfig.getItems(IDataxGlobalCfg.KEY_DISPLAY_NAME));
429429
}
430430

431+
@Override
432+
public EndType getEndType() {
433+
return EndType.Workflow;
434+
}
435+
431436
public boolean validateName(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
432437
UploadPluginMeta pluginMeta = (UploadPluginMeta) context.get(UploadPluginMeta.KEY_PLUGIN_META);
433438
Objects.requireNonNull(pluginMeta, "pluginMeta can not be null");

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

Lines changed: 11 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -24,22 +24,20 @@
2424
import com.qlangtech.tis.datax.DataXName;
2525
import com.qlangtech.tis.datax.DefaultDataXProcessorManipulate;
2626
import com.qlangtech.tis.datax.IDataxGlobalCfg;
27-
import com.qlangtech.tis.datax.IDataxProcessor;
2827
import com.qlangtech.tis.datax.IDataxReader;
28+
import com.qlangtech.tis.datax.StoreResourceType;
2929
import com.qlangtech.tis.datax.StoreResourceTypeConstants;
3030
import com.qlangtech.tis.datax.impl.DataxProcessor;
3131
import com.qlangtech.tis.datax.impl.TransformerInfo;
3232
import com.qlangtech.tis.extension.Descriptor;
3333
import com.qlangtech.tis.extension.IDescribableManipulate;
3434
import com.qlangtech.tis.extension.TISExtension;
35-
import com.qlangtech.tis.extension.impl.XmlFile;
3635
import com.qlangtech.tis.manage.IAppSource;
3736
import com.qlangtech.tis.manage.biz.dal.pojo.AppType;
3837
import com.qlangtech.tis.manage.biz.dal.pojo.Application;
3938
import com.qlangtech.tis.manage.common.AppAndRuntime;
39+
import com.qlangtech.tis.plugin.IEndTypeGetter;
4040
import com.qlangtech.tis.plugin.IPluginStore;
41-
import com.qlangtech.tis.datax.StoreResourceType;
42-
import com.qlangtech.tis.plugin.KeyedPluginStore.Key;
4341
import com.qlangtech.tis.plugin.annotation.FormField;
4442
import com.qlangtech.tis.plugin.annotation.FormFieldType;
4543
import com.qlangtech.tis.plugin.annotation.Validator;
@@ -49,18 +47,14 @@
4947
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
5048
import com.qlangtech.tis.sql.parser.tuple.creator.IStreamIncrGenerateStrategy;
5149
import com.qlangtech.tis.util.IPluginContext;
52-
import com.qlangtech.tis.util.TransformerRuleKey;
5350
import com.qlangtech.tis.util.UploadPluginMeta;
5451
import org.apache.commons.io.FileUtils;
55-
import org.apache.commons.io.filefilter.FalseFileFilter;
56-
import org.apache.commons.io.filefilter.SuffixFileFilter;
5752
import org.apache.commons.lang3.StringUtils;
5853
import org.apache.commons.lang3.tuple.Pair;
5954

6055
import java.io.File;
6156
import java.io.IOException;
62-
import java.util.Collection;
63-
import java.util.Collections;
57+
import java.util.Date;
6458
import java.util.HashSet;
6559
import java.util.List;
6660
import java.util.Map;
@@ -103,6 +97,9 @@ public Application buildApp() {
10397
app.setProjectName(this.name);
10498
app.setDptId(Integer.parseInt(this.dptId));
10599
app.setRecept(this.recept);
100+
app.setUpdateTime(new Date());
101+
app.setLastProcessTime(new Date());
102+
106103
app.setAppType(AppType.DataXPipe.getType());
107104
return app;
108105
}
@@ -129,36 +126,13 @@ public Pair<List<RecordTransformerRules>, IPluginStore> getRecordTransformerRule
129126
public Set<TransformerInfo> getTransformerInfo(IPluginContext pluginCtx, Map<String, List<DBDataXChildTask>> groupedChildTask) {
130127
Set<TransformerInfo> tinfos = new HashSet<>();
131128
addTransformerInfo(tinfos, pluginCtx, groupedChildTask, this.getResType(), this.identityValue(), (tableName, context) -> {
132-
// return Optional<RecordTransformerRules>
133129
Pair<List<RecordTransformerRules>, IPluginStore> tabTransformerRule
134130
= DataFlowDataXProcessor.loadRecordTransformerRulesAndPluginStore(context, this.getResType(), this.name, tableName);
135131
for (RecordTransformerRules trule : tabTransformerRule.getKey()) {
136132
return Optional.of(trule);
137133
}
138134
return Optional.empty();
139135
});
140-
141-
// Key transformerRuleKey = TransformerRuleKey.createStoreKey(
142-
// pluginCtx, this.getResType(), this.identityValue(), "dump");
143-
// XmlFile sotre = transformerRuleKey.getSotreFile();
144-
// File parent = sotre.getFile().getParentFile();
145-
// if (!parent.exists()) {
146-
// return Collections.emptySet();
147-
// }
148-
// Optional<RecordTransformerRules> transformerRules = null;
149-
// String xmlExtend = Descriptor.getPluginFileName(org.apache.commons.lang.StringUtils.EMPTY);
150-
// SuffixFileFilter filter = new SuffixFileFilter(xmlExtend);
151-
// Collection<File> matched = FileUtils.listFiles(parent, filter, FalseFileFilter.INSTANCE);
152-
// for (File tfile : matched) {
153-
// String tabName = org.apache.commons.lang.StringUtils.substringBefore(tfile.getName(), xmlExtend);
154-
// if (groupedChildTask.containsKey(tabName)) {
155-
// transformerRules = RecordTransformerRules.loadTransformerRules(
156-
// pluginCtx, this.getResType(), this.identityValue(), tabName);
157-
// if (transformerRules.isPresent()) {
158-
// tinfos.add(new TransformerInfo(tabName, transformerRules.get().rules.size()));
159-
// }
160-
// }
161-
// }
162136
return tinfos;
163137
}
164138

@@ -218,14 +192,18 @@ private <T> T writerPluginOverwrite(Function<IStreamIncrGenerateStrategy, T> fun
218192

219193

220194
@TISExtension()
221-
public static class DescriptorImpl extends Descriptor<IAppSource> implements IDescribableManipulate<DefaultDataXProcessorManipulate> {
195+
public static class DescriptorImpl extends Descriptor<IAppSource> implements IDescribableManipulate<DefaultDataXProcessorManipulate>, IEndTypeGetter {
222196
static final int MAX_RECEPT_LENGTH = 20;
223197

224198
public DescriptorImpl() {
225199
super();
226200
this.registerSelectOptions(KEY_FIELD_NAME, () -> ParamsConfig.getItems(IDataxGlobalCfg.KEY_DISPLAY_NAME));
227201
}
228202

203+
@Override
204+
public EndType getEndType() {
205+
return EndType.Pipeline;
206+
}
229207

230208
public boolean validateRecept(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
231209

@@ -237,12 +215,6 @@ public boolean validateRecept(IFieldErrorHandler msgHandler, Context context, St
237215
}
238216

239217
public boolean validateName(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
240-
241-
// if (PATTERN_START_WITH_NUMBER.matcher(value).matches()) {
242-
// msgHandler.addFieldError(context, fieldName, "不能以数字开头");
243-
// return false;
244-
// }
245-
246218
UploadPluginMeta pluginMeta = (UploadPluginMeta) context.get(UploadPluginMeta.KEY_PLUGIN_META);
247219
Objects.requireNonNull(pluginMeta, "pluginMeta can not be null");
248220
if (pluginMeta.isUpdate()) {

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,6 @@
3131
import com.qlangtech.tis.plugin.annotation.SubForm;
3232
import com.qlangtech.tis.plugin.annotation.Validator;
3333
import com.qlangtech.tis.plugin.datax.SelectedTab;
34-
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
3534
import com.qlangtech.tis.plugin.ds.CMeta;
3635
import com.qlangtech.tis.plugin.ds.ColumnMetaData;
3736
import com.qlangtech.tis.plugin.ds.DataSourceFactory;
@@ -211,7 +210,7 @@ public void startScanDependency() {
211210

212211
@Override
213212
public DS getDataSourceFactory() {
214-
return TIS.getDataBasePlugin(PostedDSProp.parse(this.dbName));
213+
return DataSourceFactory.load(this.dbName);
215214
}
216215

217216
public final List<ColumnMetaData> getTableMetadata(EntityName table) throws TableNotFoundException {

tis-datax/tis-datax-common-plugin/src/main/java/com/qlangtech/tis/plugin/datax/format/guesstype/GuessOff.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@
3232
**/
3333
public class GuessOff extends GuessFieldType {
3434
@Override
35-
public Map<String, DataType> processStructGuess(IGuessColTypeFormatConfig textFormat, StructuredReader reader) throws IOException {
35+
public Map<String, DataType> processStructGuess(TargetTabsEntities targetTabs, IGuessColTypeFormatConfig textFormat, StructuredReader reader) throws IOException {
3636
return Collections.emptyMap();
3737
}
3838

tis-datax/tis-datax-common-plugin/src/main/java/com/qlangtech/tis/plugin/datax/format/guesstype/GuessOn.java

Lines changed: 21 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,6 @@
3232
import com.qlangtech.tis.plugin.ds.DataTypeMeta;
3333
import com.qlangtech.tis.plugin.ds.JDBCTypes;
3434
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
35-
import org.apache.commons.collections.CollectionUtils;
3635
import org.apache.commons.collections.MapUtils;
3736
import org.apache.commons.lang.StringUtils;
3837
import org.slf4j.Logger;
@@ -58,16 +57,18 @@ public class GuessOn extends GuessFieldType {
5857

5958

6059
@Override
61-
public Map<String, Map<String, DataType>> processStructGuess(
60+
public Map<KafkaLogicalTableName, Map<String, DataType>>
61+
processStructGuess(
62+
TargetTabsEntities targetTabs,
6263
IGuessColTypeFormatConfig textFormat, StructuredReader<StructuredRecord> reader) throws IOException {
63-
Map<String, Map<String, PriorityDataType>> result = Maps.newHashMap();
64+
Map<KafkaLogicalTableName, Map<String, PriorityDataType>> result = Maps.newHashMap();
6465
Map<String, PriorityDataType> priorityResult = null;
6566
// priorityResult = Maps.newHashMap();
6667
int lineIndex = 0;
6768
StructuredRecord row = null;
68-
// Map<String, Object> rowVals;
69-
// PriorityDataType guessType = null;
70-
String tabName = null;
69+
70+
PhysicsTable2LogicalTableMapper p2lMapper = new PhysicsTable2LogicalTableMapper(targetTabs);
71+
KafkaLogicalTableName tabName = null;
7172
while (reader.hasNext() && lineIndex++ < maxInspectLine) {
7273
if (lineIndex % 1000 == 0) {
7374
logger.info("has scan rows:{}", lineIndex);
@@ -76,8 +77,8 @@ public Map<String, Map<String, DataType>> processStructGuess(
7677
if (row == null) {
7778
continue;
7879
}
79-
tabName = row.tabName;// StringUtils.defaultString(, DEFAUTL_TABLE_NAME);
80-
if (StringUtils.isEmpty(tabName)) {
80+
tabName = p2lMapper.parseLogicalTableName(row.tabName);
81+
if (tabName == null) {
8182
throw new IllegalStateException("tableName can not be empty");
8283
}
8384
if ((priorityResult = result.get(tabName)) == null) {
@@ -92,17 +93,21 @@ public Map<String, Map<String, DataType>> processStructGuess(
9293
}
9394
return result.entrySet().stream().collect(Collectors.toMap((e) -> e.getKey()
9495
, (e) -> {
95-
return e.getValue().entrySet().stream().collect(Collectors.toMap((col) -> col.getKey(), (col) -> {
96-
DataType type = null;
97-
if ((type = col.getValue().type) != null) {
98-
return type;
99-
} else {
100-
return defaultDataTypeForNullVal();
101-
}
102-
}));
96+
return e.getValue().entrySet().stream().collect(
97+
Collectors.toMap(
98+
(col) -> col.getKey()
99+
, (col) -> {
100+
DataType type = null;
101+
if ((type = col.getValue().type) != null) {
102+
return type;
103+
} else {
104+
return defaultDataTypeForNullVal();
105+
}
106+
}));
103107
}));
104108
}
105109

110+
106111
private void parseStructedRecordColType(
107112
IGuessColTypeFormatConfig textFormat, StructuredRecord row, Map<String, PriorityDataType> priorityResult) {
108113
Map<String, Object> rowVals;

tis-datax/tis-datax-common-plugin/src/main/resources/com/qlangtech/tis/plugin/datax/DefaultDataxProcessor.json

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,14 +6,17 @@
66
"dptId": {
77
"label": "所属部门",
88
"enum": "com.qlangtech.tis.coredefine.module.action.DataxAction.getDepartments():uncache_true",
9+
"dftVal": "com.qlangtech.tis.manage.biz.dal.pojo.Department.dftDepartmentId()",
910
"creator": {
1011
"routerLink": "/base/departmentlist",
1112
"label": "部门管理"
1213
}
1314
},
1415
"recept": {
1516
"label": "接口人",
16-
"placeholder": "小明"
17+
"placeholder": "小明",
18+
"dftVal": "com.qlangtech.tis.manage.common.UserUtils.currentLoginUserName():uncache_true",
19+
"help": "数据通道的业务联系人,如运行生命周期内有任何问题可以向他咨询"
1720
},
1821
"globalCfg": {
1922
"label": "全局配置",

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
"val": "utf8"
2626
}
2727
],
28+
"dftVal": "utf8",
2829
"help": "数据数据"
2930
},
3031
"extraParams": {

tis-datax/tis-datax-dolphinscheduler-plugin/src/main/java/com/qlangtech/tis/plugin/datax/doplinscheduler/history/DSWorkFlowBuildHistoryPayload.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.qlangtech.tis.dao.ICommonDAOContext;
2424
import com.qlangtech.tis.datax.DefaultDataXProcessorManipulate;
2525
import com.qlangtech.tis.datax.IDataxProcessor;
26+
import com.qlangtech.tis.manage.IAppSource;
2627
import com.qlangtech.tis.plugin.IPluginStore;
2728
import com.qlangtech.tis.plugin.datax.WorkFlowBuildHistoryPayload;
2829
import com.qlangtech.tis.plugin.datax.doplinscheduler.export.DolphinSchedulerURLBuilder.DolphinSchedulerResponse;
@@ -45,7 +46,7 @@ public class DSWorkFlowBuildHistoryPayload extends WorkFlowBuildHistoryPayload {
4546
public DSWorkFlowBuildHistoryPayload(IDataxProcessor dataxProcessor, Integer tisTaskId, ICommonDAOContext daoContext) {
4647
super(dataxProcessor, tisTaskId, daoContext);
4748
Pair<List<ExportTISPipelineToDolphinscheduler>, IPluginStore<DefaultDataXProcessorManipulate>> pluginStorePair
48-
= DefaultDataXProcessorManipulate.loadPlugins(null, ExportTISPipelineToDolphinscheduler.class, this.dataxProcessor.getDataXName());
49+
= DefaultDataXProcessorManipulate.loadPlugins(null, ExportTISPipelineToDolphinscheduler.class, ((IAppSource) this.dataxProcessor).getDataXName());
4950
for (ExportTISPipelineToDolphinscheduler exportDSCfg : pluginStorePair.getLeft()) {
5051
this.exportDSCfg = exportDSCfg;
5152
return;

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

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020

2121
import com.alibaba.datax.plugin.writer.doriswriter.DorisWriterKeys;
2222
import com.alibaba.fastjson.JSONObject;
23-
import com.qlangtech.tis.TIS;
2423
import com.qlangtech.tis.annotation.Public;
2524
import com.qlangtech.tis.datax.IDataxContext;
2625
import com.qlangtech.tis.datax.IDataxProcessor;
@@ -40,7 +39,6 @@
4039

4140
import java.util.Collections;
4241
import java.util.List;
43-
import java.util.Objects;
4442
import java.util.Optional;
4543
import java.util.stream.Collectors;
4644

@@ -114,16 +112,19 @@ public static String getDftTemplate() {
114112

115113
@TISExtension()
116114
public static class DefaultDescriptor extends BaseDescriptor implements DataxWriter.IRewriteSuFormProperties {
115+
private final Descriptor<SelectedTab> dorisTabDesc;
117116

118117
public DefaultDescriptor() {
119118
super();
119+
this.dorisTabDesc = new DorisSelectedTab.DefaultDescriptor();
120120
}
121121

122122
@Override
123123
public Descriptor<SelectedTab> getRewriterSelectTabDescriptor() {
124-
Class targetClass = DorisSelectedTab.class;
125-
return Objects.requireNonNull(TIS.get().getDescriptor(targetClass)
126-
, "subForm clazz:" + targetClass + " can not find relevant Descriptor");
124+
// Class targetClass = DorisSelectedTab.class;
125+
// return Objects.requireNonNull(TIS.get().getDescriptor(targetClass)
126+
// , "subForm clazz:" + targetClass + " can not find relevant Descriptor");
127+
return this.dorisTabDesc;
127128
}
128129

129130
@Override

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

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -32,16 +32,13 @@
3232
import com.qlangtech.tis.plugin.annotation.FormFieldType;
3333
import com.qlangtech.tis.plugin.annotation.Validator;
3434
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
35-
import com.qlangtech.tis.plugin.ds.ColumnMetaData;
3635
import com.qlangtech.tis.plugin.ds.DBConfig;
37-
3836
import com.qlangtech.tis.plugin.ds.DataType;
3937
import com.qlangtech.tis.plugin.ds.DataType.DefaultTypeVisitor;
4038
import com.qlangtech.tis.plugin.ds.JDBCConnection;
4139
import com.qlangtech.tis.plugin.ds.JDBCTypes;
4240
import com.qlangtech.tis.plugin.ds.NoneSplitTableStrategy;
4341
import com.qlangtech.tis.plugin.ds.SplitTableStrategy;
44-
import com.qlangtech.tis.plugin.ds.TableNotFoundException;
4542
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
4643
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
4744
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
@@ -56,6 +53,7 @@
5653
import java.net.URL;
5754
import java.sql.ResultSet;
5855
import java.sql.SQLException;
56+
import java.time.Duration;
5957
import java.util.Collections;
6058
import java.util.List;
6159
import java.util.Map;
@@ -246,19 +244,25 @@ protected boolean validateDSFactory(final IControlMsgHandler msgHandler, final C
246244
if (valid) {
247245

248246
final DorisSourceFactory dorisDS = (DorisSourceFactory) dsFactory;
249-
247+
final Duration socketReadTimeout = Duration.ofSeconds(5);
250248
for (String feLoadHost : dorisDS.getLoadUrls()) {
251249
//利用doris的clusterAction: https://doris.apache.org/zh-CN/docs/1.2/admin-manual/http-actions/fe/cluster-action
252250
// 对:{"msg":"success","code":0,"data":{"http":["192.168.28.200:8030"],"mysql":["192.168.28.200:9030"]},"count":0}
253251
StringBuffer clusterInfoApiUrl = new StringBuffer("http://");
254252
clusterInfoApiUrl.append(feLoadHost).append("/rest/v2/manager/cluster/cluster_info/conn_info");
253+
255254
try {
256255
Boolean success = HttpUtils.get(new URL(clusterInfoApiUrl.toString()), new PostFormStreamProcess<Boolean>() {
257256
@Override
258257
public ContentType getContentType() {
259258
return ContentType.JSON;
260259
}
261260

261+
@Override
262+
public Duration getSocketReadTimeout() {
263+
return socketReadTimeout;
264+
}
265+
262266
@Override
263267
public void preSet(HttpURLConnection conn) throws IOException {
264268
super.preSet(conn);

0 commit comments

Comments
 (0)