Skip to content

Commit 8767d99

Browse files
committed
remove TableAlias
1 parent 9e5e047 commit 8767d99

53 files changed

Lines changed: 277 additions & 288 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

tis-datax/executor/tis-datax-executor/launch-mysql-pipeline-on-grpc.sh

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,4 +22,4 @@ java -Ddata.dir=/opt/data/tis -Denv_props=true -Dlog.dir=/opt/logs/tis -Druntime
2222
-Dlogback.configurationFile=logback-datax.xml -DexecTimeStamp=0 \
2323
-classpath ./lib/*:./tis-datax-executor.jar:./conf/ \
2424
-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=50002 \
25-
com.qlangtech.tis.plugin.datax.DataXPipelinePreviewMain dameng_mysql 51509
25+
com.qlangtech.tis.plugin.datax.DataXPipelinePreviewMain dfs_mysql2 51509

tis-datax/executor/tis-datax-executor/launch-mysql-pipeline.sh

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,6 @@
1919
cd /opt/tis/tis-datax-executor
2020

2121
java -Ddata.dir=/opt/data/tis -Denv_props=true -Dlog.dir=/opt/logs/tis -Druntime=daily -Dlogback.configurationFile=logback-datax.xml -DexecTimeStamp=0 \
22-
-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=50002 \
22+
-agentlib:jdwp=transport=dt_socket,server=y,suspend=n,address=50002 \
2323
-classpath ./lib/*:./tis-datax-executor.jar:./conf/ \
24-
com.qlangtech.tis.plugin.datax.grpc.DefaultDataXPreviewRocrdsImpl mysql_elastic2
24+
com.qlangtech.tis.plugin.datax.grpc.DefaultDataXPreviewRocrdsImpl dfs_mysql2

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/datax/common/DataXRealExecutor.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -106,13 +106,14 @@ private boolean isNumericJdbcType(Map<String, DataType> typeMap, String colKey)
106106
* @return
107107
*/
108108
public PreviewRecords previewRecords(String tableName, QueryCriteria queryCriteria) {
109-
109+
IPluginContext pluginCtx = IPluginContext.namedContext(this.dataXName);
110110
if (StringUtils.isEmpty(tableName)) {
111111
throw new IllegalArgumentException("param tableName can not be null");
112112
}
113113
final IDataxReader dataXReader = this.getDataxReader();
114114
IGroupChildTaskIterator subTasks = dataXReader.getSubTasks((tab) -> StringUtils.equals(tab.getName(), tableName));
115-
IPluginContext pluginCtx = IPluginContext.namedContext(this.dataXName);
115+
116+
116117
while (subTasks.hasNext()) {
117118
IDataxReaderContext readerContext = subTasks.next();
118119

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/plugin/datax/grpc/DefaultDataXPreviewRocrdsImpl.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@
2727
import com.alibaba.datax.common.util.Configuration;
2828
import com.google.common.collect.Lists;
2929
import com.google.common.collect.Maps;
30-
3130
import com.qlangtech.tis.TIS;
3231
import com.qlangtech.tis.datax.TISJarLoader;
3332
import com.qlangtech.tis.datax.common.DataXRealExecutor;
@@ -48,7 +47,6 @@
4847
import org.slf4j.LoggerFactory;
4948

5049
import java.util.Arrays;
51-
import java.util.Collections;
5250
import java.util.List;
5351
import java.util.Map;
5452
import java.util.Objects;
@@ -85,7 +83,7 @@ public static void main(String[] args) throws Exception {
8583

8684

8785
PreviewRecords records
88-
= previewRocrds.pipeSynchronize.previewRecords("orderdetail", queryCriteria);
86+
= previewRocrds.pipeSynchronize.previewRecords("totalpayinfo", queryCriteria);
8987
records.getPageRows().forEach((r) -> {
9088
// System.out.println(r.getColumn("order_id"));
9189
System.out.println(r);

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

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,6 @@
3232
import com.qlangtech.tis.datax.IGroupChildTaskIterator;
3333
import com.qlangtech.tis.datax.StoreResourceType;
3434
import com.qlangtech.tis.datax.StoreResourceTypeConstants;
35-
import com.qlangtech.tis.datax.TableAliasMapper;
3635
import com.qlangtech.tis.datax.impl.DataXCfgGenerator;
3736
import com.qlangtech.tis.datax.impl.DataxProcessor;
3837
import com.qlangtech.tis.datax.impl.DataxReader;
@@ -415,10 +414,7 @@ public DataXCfgGenerator.GenerateCfgs getDataxCfgFileNames(IPluginContext plugin
415414
return DataxProcessor.getDataxCfgFileNames(pluginCtx, partialTrigger, this);
416415
}
417416

418-
@Override
419-
public TableAliasMapper getTabAlias(IPluginContext pluginCtx, boolean withDft) {
420-
return TableAliasMapper.Null;
421-
}
417+
422418

423419

424420
@TISExtension()

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

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,7 @@
3333
import com.qlangtech.tis.datax.IDataxReader;
3434
import com.qlangtech.tis.datax.IDataxWriter;
3535
import com.qlangtech.tis.datax.SourceColMetaGetter;
36-
import com.qlangtech.tis.datax.TableAlias;
37-
import com.qlangtech.tis.datax.TableAliasMapper;
36+
import com.qlangtech.tis.datax.StoreResourceType;
3837
import com.qlangtech.tis.datax.impl.DataxProcessor;
3938
import com.qlangtech.tis.datax.impl.DataxWriter;
4039
import com.qlangtech.tis.exec.ExecutePhaseRange;
@@ -43,7 +42,6 @@
4342
import com.qlangtech.tis.fullbuild.indexbuild.IRemoteTaskPreviousTrigger;
4443
import com.qlangtech.tis.manage.common.TisUTF8;
4544
import com.qlangtech.tis.plugin.KeyedPluginStore;
46-
import com.qlangtech.tis.datax.StoreResourceType;
4745
import com.qlangtech.tis.plugin.annotation.FormField;
4846
import com.qlangtech.tis.plugin.annotation.FormFieldType;
4947
import com.qlangtech.tis.plugin.annotation.Validator;
@@ -146,12 +144,15 @@ private class PreAndPostSQLExecutor implements IRemoteTaskPostTrigger, IRemoteTa
146144
private final IExecChainContext execContext;
147145
private final EntityName entity;
148146
private final ISelectedTab tab;
147+
private final Optional<AutoCreateTable> writerTableExecutor;
149148

150149
public PreAndPostSQLExecutor(boolean preExecute, IExecChainContext execContext, EntityName entity, ISelectedTab tab) {
151150
this.preExecute = preExecute;
152151
this.execContext = execContext;
153152
this.entity = entity;
154153
this.tab = tab;
154+
IDataxWriter writer = execContext.getProcessor().getWriter(null);
155+
this.writerTableExecutor = writer.getWriterTableExecutor();
155156
}
156157

157158
@Override
@@ -169,14 +170,17 @@ private String validateSQL(String sql) {
169170
@Override
170171
public void run() {
171172

172-
final TableAliasMapper tableAliasMapper
173-
= execContext.getAttribute(TableAlias.class.getSimpleName(), () -> {
174-
return execContext.getProcessor().getTabAlias(null, true);
175-
});
173+
// final TableAliasMapper tableAliasMapper
174+
// = execContext.getAttribute(TableAlias.class.getSimpleName(), () -> {
175+
// return execContext.getProcessor().getTabAlias(null, true);
176+
// });
177+
178+
179+
;
176180

177181
BasicDataSourceFactory dsFactory = ((BasicDataSourceFactory) getDataSourceFactory());
178182
dsFactory.visitAllConnection((conn) -> {
179-
SelectTable toTable = SelectTable.create(tableAliasMapper.get(tab).getTo(), dsFactory);
183+
SelectTable toTable = SelectTable.create(new TableMap(writerTableExecutor, tab).getTo(), dsFactory);
180184
String preSqlStatement = StringUtils.replace(validateSQL(this.preExecute ? preSql : postSql), TABLE_NAME_PLACEHOLDER, toTable.getTabName());
181185
String checkTabExist = "select 1 from " + toTable.getTabName();
182186
final AtomicBoolean tabExist = new AtomicBoolean(false);

tis-datax/tis-datax-dfs-plugin/src/main/java/com/qlangtech/tis/plugin/datax/AbstractDFSReader.java

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -22,14 +22,16 @@
2222
import com.qlangtech.tis.datax.IDataxProcessor.TableMap;
2323
import com.qlangtech.tis.datax.IGroupChildTaskIterator;
2424
import com.qlangtech.tis.datax.impl.DataxReader;
25-
import com.qlangtech.tis.extension.impl.SuFormProperties;
2625
import com.qlangtech.tis.plugin.KeyedPluginStore;
2726
import com.qlangtech.tis.plugin.annotation.FormField;
2827
import com.qlangtech.tis.plugin.annotation.FormFieldType;
2928
import com.qlangtech.tis.plugin.annotation.SubForm;
3029
import com.qlangtech.tis.plugin.annotation.Validator;
31-
import com.qlangtech.tis.plugin.datax.resmatcher.WildcardDFSResMatcher;
32-
import com.qlangtech.tis.plugin.ds.*;
30+
import com.qlangtech.tis.plugin.ds.ColumnMetaData;
31+
import com.qlangtech.tis.plugin.ds.DBIdentity;
32+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
33+
import com.qlangtech.tis.plugin.ds.TableInDB;
34+
import com.qlangtech.tis.plugin.ds.TableNotFoundException;
3335
import com.qlangtech.tis.plugin.tdfs.DFSResMatcher;
3436
import com.qlangtech.tis.plugin.tdfs.IDFSReader;
3537
import com.qlangtech.tis.plugin.tdfs.ITDFSSession;
@@ -43,10 +45,8 @@
4345
import java.util.Collections;
4446
import java.util.List;
4547
import java.util.Objects;
46-
import java.util.Optional;
4748
import java.util.function.Predicate;
4849
import java.util.function.Supplier;
49-
import java.util.stream.Collectors;
5050

5151
/**
5252
* @author: baisui 百岁
@@ -75,6 +75,11 @@ public void startScanDependency() {
7575
}
7676
}
7777

78+
// @Override
79+
public boolean isRDBMSSupport() {
80+
return Objects.requireNonNull(resMatcher, "resMatcher can not be null").isRDBMSSupport();
81+
}
82+
7883
/**
7984
* ================================================================================
8085
* support rdbms start
@@ -85,9 +90,18 @@ public void startScanDependency() {
8590

8691
public abstract List<DataXDFSReaderWithMeta.TargetResMeta> getSelectedEntities();
8792

93+
@Override
94+
public List<ISelectedTab> getSelectedTabs() {
95+
return this.resMatcher.getSelectedTabs(this);
96+
}
97+
8898
@Override
8999
public <T extends ISelectedTab> List<T> getUnfilledSelectedTabs() {
100+
//if (this.isRDBMSSupport()) {
90101
return (List<T>) selectedTabs;
102+
// } else {
103+
// return (List<T>) getSelectedTabs();
104+
// }
91105
}
92106

93107
@Override
@@ -134,11 +148,6 @@ public boolean hasMulitTable() {
134148
}
135149

136150

137-
@Override
138-
public List<ISelectedTab> getSelectedTabs() {
139-
return this.resMatcher.getSelectedTabs(this);
140-
}
141-
142151
/**
143152
* ================================================================================
144153
* support rdbms END

tis-datax/tis-datax-dfs-plugin/src/main/java/com/qlangtech/tis/plugin/datax/DataXDFSReader.java

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@
4343

4444
import java.util.Collections;
4545
import java.util.List;
46-
import java.util.Objects;
4746
import java.util.Optional;
4847
import java.util.Set;
4948
import java.util.stream.Collectors;
@@ -52,7 +51,7 @@
5251
* @author: 百岁(baisui@qlangtech.com)
5352
* @create: 2023-08-19 00:10
5453
**/
55-
public class DataXDFSReader extends AbstractDFSReader implements DataXBasicProcessMeta.IRDBMSSupport {
54+
public class DataXDFSReader extends AbstractDFSReader implements DataXBasicProcessMeta.IRDBMSSupport {
5655

5756

5857
@FormField(ordinal = 8, validate = {Validator.require})
@@ -71,7 +70,10 @@ public static List<? extends Descriptor> dfsLinkerFilter(List<? extends Descript
7170
public static List<? extends Descriptor> supportedReaderFormat(List<? extends Descriptor> descs) {
7271
return BasicPainFormatDescriptor.supportedFormat(true, descs);
7372
}
74-
73+
@Override
74+
public boolean isRDBMSSupport() {
75+
return super.isRDBMSSupport();
76+
}
7577

7678
@Override
7779
public List<DataXDFSReaderWithMeta.TargetResMeta> getSelectedEntities() {
@@ -145,10 +147,7 @@ public static String getDftTemplate() {
145147
return IOUtils.loadResourceFromClasspath(AbstractDFSReader.class, "DataXDFSReader-tpl.json");
146148
}
147149

148-
@Override
149-
public boolean isRDBMSSupport() {
150-
return Objects.requireNonNull(resMatcher, "resMatcher can not be null").isRDBMSSupport();
151-
}
150+
152151

153152
@TISExtension()
154153
public static class DefaultDescriptor extends BaseDataxReaderDescriptor {

tis-datax/tis-datax-dfs-plugin/src/main/java/com/qlangtech/tis/plugin/datax/resmatcher/BasicDFSResMatcher.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,8 @@
3939
/**
4040
* @author: 百岁(baisui@qlangtech.com)
4141
* @create: 2023-08-13 22:49
42+
* @see WildcardDFSResMatcher
43+
* @see MetaAwareDFSResMatcher
4244
**/
4345
public abstract class BasicDFSResMatcher extends DFSResMatcher {
4446

tis-datax/tis-datax-dfs-plugin/src/main/java/com/qlangtech/tis/plugin/datax/resmatcher/WildcardDFSResMatcher.java

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -65,9 +65,10 @@ public static Optional<TableMap> getTableMap(IPluginContext pluginContext, Strin
6565
throw new IllegalArgumentException("param dataXName can not be empty");
6666
}
6767
IDataxProcessor dataxProcessor = DataxProcessor.load(pluginContext, dataXName);
68-
TableAliasMapper tabAlias = dataxProcessor.getTabAlias(pluginContext, true);
69-
Optional<TableMap> tabAlia = tabAlias.getFirstTableMap();
70-
return tabAlia;
68+
return dataxProcessor.getFirstTableMap(pluginContext);
69+
// TableAliasMapper tabAlias = dataxProcessor.getTabAlias(pluginContext, true);
70+
// Optional<TableMap> tabAlia = tabAlias.getFirstTableMap();
71+
// return tabAlia;
7172
}
7273

7374

@@ -85,10 +86,10 @@ public List<ColumnMetaData> getTableMetadata(IPluginContext pluginContext, Strin
8586
*/
8687
@Override
8788
public SourceColsMeta getSourceColsMeta(ITDFSSession hdfsSession, Optional<String> entityName, String path, IDataxProcessor processor) {
88-
TableAliasMapper tabAlias = processor.getTabAlias(null, false);
89-
Optional<TableAlias> findMapper = tabAlias.findFirst();
89+
// TableAliasMapper tabAlias = processor.getTabAlias(null, false);
90+
Optional<TableMap> findMapper = processor.getFirstTableMap(null); // tabAlias.findFirst();
9091
IDataxProcessor.TableMap tabMapper
91-
= (IDataxProcessor.TableMap) findMapper.orElseThrow(() -> new NullPointerException("TableAlias can not be null"));
92+
= findMapper.orElseThrow(() -> new NullPointerException("TableAlias can not be null"));
9293
return new SourceColsMeta(tabMapper.getSourceCols());
9394
}
9495

@@ -101,8 +102,8 @@ public List<ISelectedTab> getSelectedTabs(IDFSReader dfsReader) {
101102
}
102103

103104
IDataxProcessor processor = DataxProcessor.load(IPluginContext.getThreadLocalInstance(), reader.dataXName);
104-
TableAliasMapper tabAlias = processor.getTabAlias(IPluginContext.getThreadLocalInstance(), false);
105-
Optional<TableAlias> findMapper = tabAlias.findFirst();
105+
// TableAliasMapper tabAlias = processor.getTabAlias(IPluginContext.getThreadLocalInstance(), false);
106+
Optional<TableMap> findMapper = processor.getFirstTableMap(IPluginContext.getThreadLocalInstance(),false); //tabAlias.findFirst();
106107
if (findMapper.isPresent()) {
107108
IDataxProcessor.TableMap tabMapper = (IDataxProcessor.TableMap) findMapper.get();
108109
return Collections.singletonList(

0 commit comments

Comments
 (0)