Skip to content

Commit 91c8984

Browse files
committed
com.qlangtech.tis.config.hive.meta.IHiveMetaStore.unwrapClient()
1 parent 32b1d29 commit 91c8984

9 files changed

Lines changed: 72 additions & 68 deletions

File tree

pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@
6868
<module>tis-transformer</module>
6969
<module>tis-datax/tis-datax-dolphinscheduler-plugin</module>
7070
<module>tis-datax/tis-hive-shim-common</module>
71-
<module>tis-datax/tis-datax-paimon-plugin</module>
71+
7272

7373

7474
<!-- <module>tis-solr-plugin</module>

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

Lines changed: 0 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -87,52 +87,6 @@ public Integer getRowFetchSize() {
8787
return this.fetchSize;
8888
}
8989

90-
// @Override
91-
// public Map<String, ContextParamConfig> getDBContextParams() {
92-
// return ContextParamConfig.defaultContextParams();
93-
// ContextParamConfig dbName = new ContextParamConfig("dbName") {
94-
// @Override
95-
// public ContextParamValGetter<RdbmsRunningContext> valGetter() {
96-
// return new DbNameContextParamValGetter();
97-
// }
98-
//
99-
// @Override
100-
// public DataType getDataType() {
101-
// return DataType.createVarChar(50);
102-
// }
103-
// };
104-
//
105-
// ContextParamConfig sysTimestamp = new ContextParamConfig("timestamp") {
106-
// @Override
107-
// public ContextParamValGetter<RdbmsRunningContext> valGetter() {
108-
// return new SystemTimeStampContextParamValGetter();
109-
// }
110-
//
111-
// @Override
112-
// public DataType getDataType() {
113-
// return DataType.getType(JDBCTypes.TIMESTAMP);
114-
// }
115-
// };
116-
//
117-
// ContextParamConfig tableName = new ContextParamConfig("tableName") {
118-
// @Override
119-
// public ContextParamValGetter<RdbmsRunningContext> valGetter() {
120-
// return new TableNameContextParamValGetter();
121-
// }
122-
//
123-
// @Override
124-
// public DataType getDataType() {
125-
// return DataType.createVarChar(50);
126-
// }
127-
// };
128-
//
129-
// return Lists.newArrayList(dbName, tableName, sysTimestamp)
130-
// .stream().collect(Collectors.toMap((cfg) -> cfg.getKeyName(), (cfg) -> cfg));
131-
//
132-
//// dbContextParams.put(dbName.getKeyName(), dbName);
133-
//// return dbContextParams;
134-
// }
135-
13690
@Override
13791
public final void afterSaved(IPluginContext pluginContext, Optional<Context> context) {
13892
this.preSelectedTabsHash = -1;

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -472,11 +472,16 @@ public TISDataXJobContainer(IDataXNameAware dataXName, Configuration configurati
472472
DataXJobInfo jobName) {
473473
super(configuration);
474474
this.jobArgs = args;
475-
this.jobId = args.jobId;
475+
this.jobId = Objects.requireNonNull(args.jobId);
476476
this.jobName = jobName;
477477
this.dataXName = dataXName;
478478
}
479479

480+
@Override
481+
public Integer getTaskId() {
482+
return this.jobId;
483+
}
484+
480485
@Override
481486
public String getFormatTime(TimeFormat format) {
482487
return format.format(jobArgs.getExecEpochMilli());

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

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -190,6 +190,10 @@ public static void realExecuteDump(final DataXCfgJson writerJson, IDataXPluginMe
190190
}
191191

192192
public static void realExecuteDump(String dataXName, final DataXCfgJson writerJson, IDataXPluginMeta dataxWriter) throws IllegalAccessException {
193+
realExecuteDump(null, dataXName, writerJson, dataxWriter);
194+
}
195+
196+
public static void realExecuteDump(Integer taskId, String dataXName, final DataXCfgJson writerJson, IDataXPluginMeta dataxWriter) throws IllegalAccessException {
193197
PerfTrace.getInstance(false, -1111, -1111, 0, false);
194198
final IReaderPluginMeta readerMeta = new IReaderPluginMeta() {
195199
@Override
@@ -268,22 +272,27 @@ public Configuration getWriterJsonCfg() {
268272
}
269273
};
270274

271-
realExecuteDump(dataXName, readerMeta, writerMeta);
275+
realExecuteDump(taskId, dataXName, readerMeta, writerMeta);
272276
}
273277

274278
public static void realExecuteDump(IReaderPluginMeta readerPluginMeta, IWriterPluginMeta writerMeta
275279
) throws IllegalAccessException {
276280
throw new UnsupportedOperationException("please use realExecuteDump(final String dataXName, IReaderPluginMeta readerPluginMeta, IWriterPluginMeta writerMeta");
277281
}
278282

283+
public static void realExecuteDump(final String dataXName, IReaderPluginMeta readerPluginMeta, IWriterPluginMeta writerMeta
284+
) throws IllegalAccessException {
285+
realExecuteDump(null, dataXName, readerPluginMeta, writerMeta);
286+
}
287+
279288
/**
280289
* dataXWriter执行
281290
*
282291
* @param
283292
* @param writerMeta
284293
* @throws IllegalAccessException
285294
*/
286-
public static void realExecuteDump(final String dataXName, IReaderPluginMeta readerPluginMeta, IWriterPluginMeta writerMeta
295+
public static void realExecuteDump(Integer taskId, final String dataXName, IReaderPluginMeta readerPluginMeta, IWriterPluginMeta writerMeta
287296
) throws IllegalAccessException {
288297
// final JarLoader uberClassLoader = new JarLoader(new String[]{"."});
289298
final JarLoader uberClassLoader = new TISJarLoader(TIS.get().getPluginManager());
@@ -326,6 +335,11 @@ public static void realExecuteDump(final String dataXName, IReaderPluginMeta rea
326335
LoadUtil.bind(allConf);
327336

328337
JobContainer container = new JobContainer(allConf) {
338+
@Override
339+
public Integer getTaskId() {
340+
return taskId != null ? taskId : super.getTaskId();
341+
}
342+
329343
@Override
330344
public int getTaskSerializeNum() {
331345
return super.getTaskSerializeNum();

tis-datax/tis-hive-flat-table-builder-plugin/src/main/java/com/qlangtech/tis/hive/DefaultHiveConnGetter.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -221,7 +221,7 @@ public IHiveMetaStore visit(IKerberosUserToken token) {
221221
private IHiveMetaStore createHiveMetaStore() {
222222
try {
223223
final IMetaStoreClient storeClient = Hive.getWithFastCheck(hiveCfg, false).getMSC();
224-
return new DefaultHiveMetaStore(storeClient, metaStoreUrls);
224+
return new DefaultHiveMetaStore(hiveCfg, storeClient, metaStoreUrls);
225225
} catch (Exception e) {
226226
// if (ExceptionUtils.indexOfThrowable(e, java.net.ConnectException.class) > -1) {
227227
// throw TisException.create(metaStoreUrls, e);

tis-datax/tis-hive-flat-table-builder-plugin/src/main/java/com/qlangtech/tis/hive/DefaultHiveMetaStore.java

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
import com.qlangtech.tis.config.hive.meta.IHiveMetaStore;
77
import com.qlangtech.tis.config.hive.meta.PartitionFilter;
88
import com.qlangtech.tis.hive.shim.IHiveSerDe;
9+
import org.apache.hadoop.hive.conf.HiveConf;
910
import org.apache.hadoop.hive.metastore.IMetaStoreClient;
1011
import org.apache.hadoop.hive.metastore.api.NoSuchObjectException;
1112
import org.apache.hadoop.hive.metastore.api.SerDeInfo;
@@ -39,11 +40,13 @@
3940
public class DefaultHiveMetaStore implements IHiveMetaStore {
4041
final IMetaStoreClient storeClient;
4142
private final String metaStoreUrls;
43+
private final HiveConf hiveCfg;
4244
private static final Logger logger = LoggerFactory.getLogger(DefaultHiveMetaStore.class);
4345

44-
public DefaultHiveMetaStore(IMetaStoreClient storeClient, String metaStoreUrls) {
46+
public DefaultHiveMetaStore(HiveConf hiveCfg, IMetaStoreClient storeClient, String metaStoreUrls) {
4547
this.storeClient = storeClient;
4648
this.metaStoreUrls = metaStoreUrls;
49+
this.hiveCfg = hiveCfg;
4750
}
4851

4952
@Override
@@ -108,6 +111,16 @@ public String getStorageLocation() {
108111
}
109112
}
110113

114+
@Override
115+
public HiveConf getHiveCfg() {
116+
return Objects.requireNonNull(this.hiveCfg, "hiveCfg can not be null");
117+
}
118+
119+
@Override
120+
public IMetaStoreClient unwrapClient() {
121+
return this.storeClient;
122+
}
123+
111124
public static class HiveStoredAs extends StoredAs {
112125
private final SerDeInfo serdeInfo;
113126
private final InputFormat inputFormat;

tis-incr/tis-flink-dependency/pom.xml

Lines changed: 0 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -35,22 +35,7 @@
3535
</properties>
3636
<dependencies>
3737

38-
<dependency>
39-
<groupId>org.apache.flink</groupId>
40-
<artifactId>flink-cdc-common</artifactId>
41-
<version>${flink.cdc.version}</version>
42-
</dependency>
4338

44-
<dependency>
45-
<groupId>org.apache.flink</groupId>
46-
<artifactId>flink-cdc-runtime</artifactId>
47-
<version>${flink.cdc.version}</version>
48-
</dependency>
49-
<dependency>
50-
<groupId>org.apache.flink</groupId>
51-
<artifactId>flink-cdc-composer-tis</artifactId>
52-
<version>${flink.cdc.version}</version>
53-
</dependency>
5439
<!-- <dependency>-->
5540
<!-- <groupId>com.qlangtech.tis</groupId>-->
5641
<!-- <artifactId>tis-dag</artifactId>-->

tis-incr/tis-flink-extends/pom.xml

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,29 @@
3535

3636
<dependencies>
3737

38+
<dependency>
39+
<groupId>org.apache.flink</groupId>
40+
<artifactId>flink-cdc-common</artifactId>
41+
<version>${flink.cdc.version}</version>
42+
</dependency>
43+
44+
<dependency>
45+
<groupId>org.apache.flink</groupId>
46+
<artifactId>flink-cdc-runtime</artifactId>
47+
<version>${flink.cdc.version}</version>
48+
<exclusions>
49+
<exclusion>
50+
<groupId>org.apache.calcite</groupId>
51+
<artifactId>calcite-core</artifactId>
52+
</exclusion>
53+
</exclusions>
54+
</dependency>
55+
<dependency>
56+
<groupId>org.apache.flink</groupId>
57+
<artifactId>flink-cdc-composer-tis</artifactId>
58+
<version>${flink.cdc.version}</version>
59+
</dependency>
60+
3861
<dependency>
3962
<groupId>org.apache.flink</groupId>
4063
<artifactId>flink-connector-jdbc</artifactId>
@@ -144,6 +167,16 @@
144167
<exclude>org.codehaus.groovy:groovy-all</exclude>
145168
</excludes>
146169
</artifactSet>
170+
<filters>
171+
<filter>
172+
<artifact>*:*</artifact>
173+
<excludes>
174+
<exclude>META-INF/*.SF</exclude>
175+
<exclude>META-INF/*.DSA</exclude>
176+
<exclude>META-INF/*.RSA</exclude>
177+
</excludes>
178+
</filter>
179+
</filters>
147180
<!-- <transformers>-->
148181
<!-- <transformer implementation="org.apache.maven.plugins.shade.resource.AppendingTransformer">-->
149182
<!-- <resource>reference.conf</resource>-->

tis-incr/tis-flink-extends/scp.sh

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
#scp ./target/tis-flink-extends-dist-3.6.0.jar root@192.168.28.201:/tmp/release/tis/flink/lib/
22
#scp ./target/tis-flink-extends-dist-4.1.0-SNAPSHOT.jar root@192.168.28.201:/tmp/release/tis/flink/lib/
33

4-
scp ./target/tis-flink-extends-dist-4.1.0-SNAPSHOT.jar root@192.168.28.200:/tmp/flink/lib/
4+
scp ./target/tis-flink-extends-dist-4.3.0-SNAPSHOT.jar root@192.168.28.201:/tmp/release/tis/flink/lib/
55

0 commit comments

Comments
 (0)