Skip to content

Commit 54e940a

Browse files
committed
1. 删除批量构建历史的任务https://github.com/datavane/tis/issues/487,2. 添加基于TIS transformer的主表与维表的宽表构建:datavane/tis#483
1 parent 8767d99 commit 54e940a

18 files changed

Lines changed: 497 additions & 52 deletions

File tree

pom.xml

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -685,6 +685,15 @@
685685
<version>${project.version}</version>
686686
</dependency>
687687

688+
<!-- 解决IDE对jdk.tools依赖的报错,该依赖实际上会被exclusion排除 -->
689+
<dependency>
690+
<groupId>jdk.tools</groupId>
691+
<artifactId>jdk.tools</artifactId>
692+
<version>1.8</version>
693+
<scope>system</scope>
694+
<systemPath>${java.home}/../lib/tools.jar</systemPath>
695+
</dependency>
696+
688697
</dependencies>
689698
</dependencyManagement>
690699

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -248,8 +248,6 @@ public Optional<IPluginStore<DefaultDataXProcessorManipulate>> getManipulateStor
248248
}
249249

250250
return Optional.of(DefaultDataXProcessorManipulate.getPluginStore(null, appName));
251-
252-
// return Optional.of(DefaultDataXProcessorManipulate.loadPlugins(null, DefaultDataXProcessorManipulate.class, appName).getValue());
253251
}
254252
}
255253

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,6 @@
3434
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
3535
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
3636
import com.qlangtech.tis.util.IPluginContext;
37-
//import com.qlangtech.tis.zeppelin.TISZeppelinClient;
3837
import org.apache.commons.io.IOUtils;
3938
import org.apache.commons.lang.StringUtils;
4039
import org.slf4j.Logger;
@@ -45,7 +44,6 @@
4544
import java.sql.ResultSet;
4645
import java.sql.SQLException;
4746
import java.sql.Statement;
48-
import java.time.ZoneId;
4947
import java.util.ArrayList;
5048
import java.util.Arrays;
5149
import java.util.List;

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.qlangtech.tis.extension.Descriptor;
2424
import com.qlangtech.tis.extension.impl.IOUtils;
2525
import com.qlangtech.tis.trigger.util.JsonUtil;
26+
import com.qlangtech.tis.util.DefaultDescriptorsJSON;
2627
import com.qlangtech.tis.util.DescriptorsJSON;
2728
import junit.framework.TestCase;
2829

@@ -42,7 +43,8 @@ public void testDescriptJSONGenerate() {
4243
assertNotNull(descriptor);
4344

4445
List<Descriptor<ParamsConfig>> singleton = Collections.singletonList(descriptor);
45-
DescriptorsJSON descriptorsJSON = new DescriptorsJSON(singleton);
46+
@SuppressWarnings("all")
47+
DescriptorsJSON descriptorsJSON = new DefaultDescriptorsJSON(singleton);
4648

4749
JSON.parseObject(IOUtils.loadResourceFromClasspath(TestDataXGlobalConfig.class, "dataXGlobalConfig-descriptor-assert.json"));
4850

tis-datax/tis-datax-elasticsearch-plugin/src/main/java/com/qlangtech/tis/plugin/datax/DataXElasticsearchWriter.java

Lines changed: 46 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -135,11 +135,11 @@ public ElasticEndpoint getToken() {
135135
}
136136

137137
@Override
138-
public boolean hasDifferWithSource(IPluginContext pluginCtx, ISelectedTab esTab, IDataxProcessor.TableMap tableAlias) {
138+
public boolean hasDifferWithSource(IPluginContext pluginCtx, ESTableAlias esTab) {
139139
List<IColMetaGetter> cols = esTab.overwriteCols(pluginCtx, false);
140140
// IColMetaGetter col = null;
141141
// ISchemaField schemaCol = null;
142-
ISchema schema = convert2Schema(tableAlias);
142+
ISchema schema = convert2Schema(esTab);
143143
List<ISchemaField> schemaFields = schema.getSchemaFields();
144144
if (schemaFields.size() != cols.size()) {
145145
return true;
@@ -242,31 +242,38 @@ public Void varcharType(DataType type) {
242242

243243
@Override
244244
public List<ESColumn> initialIndex(IDataxProcessor dataxProcessor) {
245-
ESTableAlias esSchema = null;
246245

247-
Optional<IDataxProcessor.TableMap> first = dataxProcessor.getFirstTableMap(null);
248246

249-
// Optional<TableAlias> first = dataxProcessor.getTabAlias(null, true).findFirst();
250-
if (first.isPresent()) {
251-
IDataxProcessor.TableMap value = first.get();
252-
if (!(value instanceof ESTableAlias)) {
253-
throw new IllegalStateException("value must be type of 'ESTableAlias',but now is :" + value.getClass());
254-
}
255-
esSchema = (ESTableAlias) value;
256-
}
247+
// Optional<IDataxProcessor.TableMap> first = dataxProcessor.getFirstTableMap(null);
248+
ESTableAlias esSchema = getEsTableAlias(dataxProcessor);
257249

258250
Objects.requireNonNull(esSchema, "esSchema can not be null");
259-
List<CMeta> cols = esSchema.getSourceCols();
251+
List<CMeta> cols = esSchema.getCols();
260252
if (CollectionUtils.isEmpty(cols)) {
261253
throw new IllegalStateException("cols can not be null");
262254
}
263-
Optional<CMeta> firstPK = cols.stream().filter((c) -> c.isPk()).findFirst();
264-
if (!firstPK.isPresent()) {
255+
long pkCount = esSchema.getPrimaryKeys().size();// cols.stream().filter(CMeta::isPk).findFirst();
256+
if (pkCount < 1) {
265257
throw new IllegalStateException("has not set PK col");
266258
}
267259
return this.initialIndex(esSchema);
268260
}
269261

262+
public static ESTableAlias getEsTableAlias(IDataxProcessor dataxProcessor) {
263+
return getEsTableAlias(dataxProcessor.identityValue(), dataxProcessor.getFirstTableMap(null));
264+
}
265+
266+
private static ESTableAlias getEsTableAlias(String pipelineName, Optional<IDataxProcessor.TableMap> first) {
267+
268+
IDataxProcessor.TableMap tableMap = first.orElseThrow(() -> new IllegalStateException("index:" + pipelineName + " relevant tableMap can not be null"));
269+
ESTableAlias esSchema = null;
270+
ISelectedTab sourceTab = tableMap.getSourceTab();
271+
if (sourceTab instanceof ESTableAlias) {
272+
esSchema = (ESTableAlias) sourceTab;
273+
}
274+
return esSchema;
275+
}
276+
270277
/**
271278
* 当增量开始执行前,先需要初始化一下索引实例
272279
*
@@ -303,14 +310,12 @@ public List<ESColumn> initialIndex(ESTableAlias esSchema) {
303310
}
304311

305312
@Override
306-
public ISchema projectionFromExpertModel(IPluginContext context, ISelectedTab esTab
307-
, IDataxProcessor.TableMap tableAlias, Consumer<byte[]> schemaContentConsumer) {
308-
schemaContentConsumer.accept(((ESTableAlias) tableAlias).getSchemaByteContent());
313+
public ISchema projectionFromExpertModel(IPluginContext context, ESTableAlias tableAlias, Consumer<byte[]> schemaContentConsumer) {
314+
schemaContentConsumer.accept((tableAlias).getSchemaByteContent());
309315
return convert2Schema(tableAlias);
310316
}
311317

312-
private ISchema convert2Schema(IDataxProcessor.TableMap tableAlias) {
313-
ESTableAlias esTable = (ESTableAlias) tableAlias;
318+
private ISchema convert2Schema(ESTableAlias esTable) {
314319

315320
JSONObject body = new JSONObject();
316321
body.put("content", esTable.getSchemaContent());
@@ -391,6 +396,13 @@ public ISchema projectionFromExpertModel(JSONObject body, Predicate<ISchemaField
391396
esField.setUniqueKey(field.getBooleanValue(ISchemaField.KEY_PK));
392397
esField.setSharedKey(field.getBooleanValue(ISchemaField.KEY_SHARE_KEY));
393398

399+
if (esField.isUniqueKey()) {
400+
schema.setUniqueKey(esField.getName());
401+
}
402+
403+
if (esField.isSharedKey()) {
404+
schema.setSharedKey(esField.getName());
405+
}
394406

395407
if (!fieldAcceptPredicate.test(esField)) {
396408
continue;
@@ -537,20 +549,19 @@ public String getTemplate() {
537549

538550
@Override
539551
public IDataxContext getSubTask(Optional<IDataxProcessor.TableMap> tableMap, Optional<RecordTransformerRules> transformerRules) {
540-
541-
if (!tableMap.isPresent()) {
542-
throw new IllegalStateException("tableMap must be present");
543-
}
544-
IDataxProcessor.TableMap mapper = tableMap.get();
545-
if (!(mapper instanceof ESTableAlias)) {
546-
throw new IllegalStateException("mapper instance must be type of " + ESTableAlias.class.getSimpleName());
547-
}
548-
return new ESContext(this, (ESTableAlias) mapper);
552+
// if (!tableMap.isPresent()) {
553+
// throw new IllegalStateException("tableMap must be present");
554+
// }
555+
// IDataxProcessor.TableMap mapper = tableMap.get();
556+
// if (!(mapper instanceof ESTableAlias)) {
557+
// throw new IllegalStateException("mapper instance must be type of " + ESTableAlias.class.getSimpleName());
558+
// }
559+
return new ESContext(this, getEsTableAlias(this.index, tableMap));
549560
}
550561

551562

552563
@TISExtension()
553-
public static class DefaultDescriptor extends BaseDataxWriterDescriptor {
564+
public static class DefaultDescriptor extends BaseDataxWriterDescriptor implements IRewriteSuFormProperties {
554565
public DefaultDescriptor() {
555566
super();
556567
this.registerSelectOptions(FIELD_ENDPOINT, () -> ParamsConfig.getItems(ElasticEndpoint.KEY_ELASTIC_SEARCH_DISPLAY_NAME));
@@ -633,5 +644,10 @@ public boolean isRdbms() {
633644
public String getDisplayName() {
634645
return DATAX_NAME;
635646
}
647+
648+
@Override
649+
public ESTableAlias.DefaultDescriptor getRewriterSelectTabDescriptor() {
650+
return ESTableAlias.desc;
651+
}
636652
}
637653
}

tis-datax/tis-datax-elasticsearch-plugin/src/test/java/com/qlangtech/tis/plugin/datax/TestDataXElasticsearchWriter.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -140,18 +140,20 @@ public Class<?> getOwnerClass() {
140140
dataXWriter.dynamic = true;
141141

142142
String esSchema = IOUtils.loadResourceFromClasspath(DataXElasticsearchWriter.class, "es-schema-content.json");
143-
ESTableAlias tableMap = new ESTableAlias(esSchema);
143+
ESTableAlias tableMap = ESTableAlias.create(Optional.empty(), esSchema);
144144

145-
// tableMap.setSchemaContent(esSchema);
145+
// tableMap.setSchemaContent(esSchema);
146146

147147

148-
WriterTemplate.valiateCfgGenerate("es-datax-writer-assert.json", dataXWriter, tableMap);
148+
WriterTemplate.valiateCfgGenerate("es-datax-writer-assert.json"
149+
, dataXWriter, new IDataxProcessor.TableMap(Optional.empty(), tableMap));
149150

150151

151152
token.authToken = null;
152153
// token.sccessKeySecret = null;
153154

154-
WriterTemplate.valiateCfgGenerate("es-datax-writer-assert-without-option.json", dataXWriter, tableMap);
155+
WriterTemplate.valiateCfgGenerate("es-datax-writer-assert-without-option.json"
156+
, dataXWriter, new IDataxProcessor.TableMap(Optional.empty(), tableMap));
155157

156158

157159
EasyMock.verify(dataxReader);

tis-incr/tis-sink-elasticsearch7-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/connector/elasticsearch7/ElasticSearchSinkFactory.java

Lines changed: 15 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -118,15 +118,15 @@ public class ElasticSearchSinkFactory extends BasicTISSinkFactory<RowData> {
118118
Objects.requireNonNull(dataXWriter, "dataXWriter can not be null");
119119
ElasticEndpoint token = dataXWriter.getToken();
120120

121-
ESTableAlias esSchema = null;
122-
Optional<IDataxProcessor.TableMap> first = dataxProcessor.getFirstTableMap(null);// dataxProcessor.getTabAlias(null, false).findFirst();
123-
if (first.isPresent()) {
124-
IDataxProcessor.TableMap value = first.get();
125-
if (!(value instanceof ESTableAlias)) {
126-
throw new IllegalStateException("value must be type of 'ESTableAlias',but now is :" + value.getClass());
127-
}
128-
esSchema = (ESTableAlias) value;
129-
}
121+
ESTableAlias esSchema = DataXElasticsearchWriter.getEsTableAlias(dataxProcessor);
122+
// Optional<IDataxProcessor.TableMap> first = dataxProcessor.getFirstTableMap(null);// dataxProcessor.getTabAlias(null, false).findFirst();
123+
// if (first.isPresent()) {
124+
// IDataxProcessor.TableMap value = first.get();
125+
// if (!(value instanceof ESTableAlias)) {
126+
// throw new IllegalStateException("value must be type of 'ESTableAlias',but now is :" + value.getClass());
127+
// }
128+
// esSchema = (ESTableAlias) value;
129+
// }
130130

131131
Objects.requireNonNull(esSchema, "esSchema can not be null");
132132
// List<CMeta> cols = esSchema.getSourceCols();
@@ -210,14 +210,16 @@ public Void visit(UsernamePassword accessKey) {
210210
});
211211
// final List<FlinkCol> sourceColsMeta = FlinkCol.getAllTabColsMeta(tab.getCols(), sourceFlinkColCreator);
212212

213-
if (!StringUtils.equals(esSchema.getFrom(), tab.getName())) {
214-
throw new IllegalStateException("esSchema.getFrom():" + esSchema.getFrom() + " must be equal with tab.getName():" + tab.getName());
213+
if (!StringUtils.equals(esSchema.getName(), tab.getName())) {
214+
throw new IllegalStateException("esSchema.getFrom():" + esSchema.getName() + " must be equal with tab.getName():" + tab.getName());
215215
}
216216
Optional<SelectedTableTransformerRules> transformerOpt
217217
= SelectedTableTransformerRules.createTransformerRules(dataxProcessor.identityValue() //, esSchema
218218
, tab, sourceFlinkColCreator);
219-
return Collections.singletonMap(esSchema
220-
, new RowDataSinkFunc(esSchema, sinkBuilder.build(), primaryKeys
219+
220+
final IDataxProcessor.TableMap esTabMap = new IDataxProcessor.TableMap(Optional.empty(), esSchema);
221+
return Collections.singletonMap(esTabMap
222+
, new RowDataSinkFunc(esTabMap, sinkBuilder.build(), primaryKeys
221223
, IPluginContext.namedContext(dataxProcessor.identityValue())
222224
, tab
223225
, sourceFlinkColCreator

tis-transformer/pom.xml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,9 @@
4545
<scope>test</scope>
4646
</dependency>
4747
</dependencies>
48+
<build>
49+
50+
</build>
4851

4952

5053
</project>
Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
package com.qlangtech.tis.plugin.datax.transformer.impl;
2+
3+
import com.alibaba.datax.common.element.ColumnAwareRecord;
4+
import com.alibaba.fastjson.JSONObject;
5+
import com.google.common.collect.Lists;
6+
import com.qlangtech.tis.extension.MultiStepsSupportHost;
7+
import com.qlangtech.tis.extension.OneStepOfMultiSteps;
8+
import com.qlangtech.tis.extension.TISExtension;
9+
import com.qlangtech.tis.plugin.datax.SelectedTab;
10+
import com.qlangtech.tis.plugin.datax.transformer.InParamer;
11+
import com.qlangtech.tis.plugin.datax.transformer.OutputParameter;
12+
import com.qlangtech.tis.plugin.datax.transformer.UDFDefinition;
13+
import com.qlangtech.tis.plugin.datax.transformer.UDFDesc;
14+
import com.qlangtech.tis.plugin.datax.transformer.impl.joiner.JoinerSelectDataSource;
15+
import com.qlangtech.tis.plugin.datax.transformer.impl.joiner.JoinerSelectTable;
16+
import com.qlangtech.tis.plugin.datax.transformer.impl.joiner.JoinerSetMatchConditionAndCols;
17+
import com.qlangtech.tis.plugin.ds.CMeta;
18+
import com.qlangtech.tis.plugin.table.join.TableJoinMatchConditionCreatorFactory;
19+
import org.apache.commons.collections.CollectionUtils;
20+
21+
import java.util.List;
22+
23+
/**
24+
*
25+
* @author 百岁 (baisui@qlangtech.com)
26+
* @date 2026/1/13
27+
*/
28+
public class JoinerUDF extends UDFDefinition {
29+
30+
31+
@Override
32+
public List<OutputParameter> outParameters() {
33+
return Lists.newArrayList();
34+
}
35+
36+
@Override
37+
public List<InParamer> inParameters() {
38+
return Lists.newArrayList();
39+
}
40+
41+
@Override
42+
public void evaluate(ColumnAwareRecord record) {
43+
44+
}
45+
46+
@Override
47+
public List<UDFDesc> getLiteria() {
48+
return Lists.newArrayList();
49+
}
50+
51+
52+
@TISExtension
53+
public static class DefaultDescriptor extends UDFDefinition.BasicUDFDesc implements MultiStepsSupportHost {
54+
public DefaultDescriptor() {
55+
super();
56+
}
57+
58+
@Override
59+
public EndType getTransformerEndType() {
60+
return EndType.Constant;
61+
}
62+
63+
64+
@Override
65+
public String getDisplayName() {
66+
return "Joiner Outer Table";
67+
}
68+
69+
@Override
70+
public List<OneStepOfMultiSteps.BasicDesc> getStepDescriptionList() {
71+
return Lists.newArrayList(new JoinerSelectDataSource.Desc()
72+
, new JoinerSelectTable.Desc()
73+
, new JoinerSetMatchConditionAndCols.Desc());
74+
}
75+
76+
@Override
77+
public void appendExternalProps(JSONObject multiStepsCfg) {
78+
List<CMeta> sourceTabCols = SelectedTab.getSelectedCols();
79+
if (CollectionUtils.isEmpty(sourceTabCols)) {
80+
throw new IllegalStateException("sourceTabCols can not be empty");
81+
}
82+
multiStepsCfg.put(TableJoinMatchConditionCreatorFactory.KEY_SOURCE_TAB_COLS, sourceTabCols);
83+
}
84+
}
85+
86+
}

0 commit comments

Comments
 (0)