Skip to content

Commit 3555f29

Browse files
committed
add paimon sink plugin support
1 parent a184166 commit 3555f29

4 files changed

Lines changed: 9 additions & 12 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -103,7 +103,7 @@ public abstract class BasicDataXRdbmsWriter<DS extends DataSourceFactory> extend
103103
@FormField(ordinal = 12, type = FormFieldType.INT_NUMBER, validate = {Validator.integer})
104104
public Integer batchSize;
105105

106-
@FormField(ordinal = 10, type = FormFieldType.ENUM, validate = {Validator.require})
106+
@FormField(ordinal = 10, validate = {Validator.require})
107107
// 目标源中是否自动创建表,这样会方便不少
108108
public AutoCreateTable autoCreateTable;
109109

@@ -228,7 +228,7 @@ public AutoCreateTable getAutoCreateTableCanNotBeNull() {
228228
* @return
229229
*/
230230
@Override
231-
public boolean isGenerateCreateDDLSwitchOff() {
231+
public final boolean isGenerateCreateDDLSwitchOff() {
232232
return !getAutoCreateTableCanNotBeNull().enabled();
233233
}
234234

tis-datax/tis-datax-odps-plugin/src/main/java/com/qlangtech/tis/plugin/datax/DataXOdpsWriter.java

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -197,10 +197,6 @@ public DataflowTask createTask(ISqlTask nodeMeta, boolean isFinalNode, IExecChai
197197
return odpsTask;
198198
}
199199

200-
@Override
201-
public boolean isGenerateCreateDDLSwitchOff() {
202-
return false;
203-
}
204200

205201
@Override
206202
public ExecuteResult startTask(ITableBuildTask dumpTask) {

tis-incr/tis-flink-extends/src/test/java/com/qlangtech/plugins/incr/flink/TestTISFlinkClassLoaderFactory.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
2323
import com.qlangtech.tis.async.message.client.consumer.impl.MQListenerFactory;
2424
import com.qlangtech.tis.coredefine.module.action.TargetResName;
25+
import com.qlangtech.tis.datax.StoreResourceTypeConstants;
2526
import com.qlangtech.tis.datax.TimeFormat;
2627
import com.qlangtech.tis.datax.impl.DataxProcessor;
2728
import com.qlangtech.tis.manage.common.CenterResource;
@@ -46,6 +47,7 @@
4647
**/
4748
public class TestTISFlinkClassLoaderFactory implements TISEasyMock {
4849

50+
4951
@Test
5052
public void testBuildServerLoaderFactory() throws Exception {
5153
BlobLibraryCacheManager.ClassLoaderFactory loaderFactory;

tis-incr/tis-realtime-flink/src/main/java/com/qlangtech/tis/plugins/incr/flink/cdc/FlinkCol2Index.java

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@
1919
package com.qlangtech.tis.plugins.incr.flink.cdc;
2020

2121
import com.alibaba.datax.common.element.ICol2Index;
22+
import com.google.common.collect.ImmutableMap;
23+
import com.google.common.collect.ImmutableMap.Builder;
2224
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2325
import org.apache.commons.collections.CollectionUtils;
2426
import org.apache.commons.lang3.tuple.Pair;
@@ -40,16 +42,13 @@ public static FlinkCol2Index create(List<FlinkCol> cols) {
4042
if (CollectionUtils.isEmpty(cols)) {
4143
throw new IllegalArgumentException("param cols can not be empty");
4244
}
43-
Map<String, Pair<Integer, FlinkCol>> col2IdxBuilder = com.google.common.collect.Maps.newHashMap();
45+
Builder<String, ICol2Index.Col> col2IdxBuilder = ImmutableMap.builder();
4446
int idx = 0;
4547
for (FlinkCol col : cols) {
46-
col2IdxBuilder.put(col.name, Pair.of(idx++, col));
48+
col2IdxBuilder.put(col.name, new ICol2Index.Col(idx++, col.colType));
4749
}
4850

49-
return new FlinkCol2Index(
50-
col2IdxBuilder.entrySet().stream().collect(
51-
Collectors.toUnmodifiableMap((entry) -> entry.getKey() //
52-
, (entry) -> new ICol2Index.Col(entry.getValue().getKey(), entry.getValue().getValue().colType))));
51+
return new FlinkCol2Index(col2IdxBuilder.build());
5352
}
5453

5554

0 commit comments

Comments
 (0)