Skip to content

Commit cf5c04c

Browse files
committed
fix bug for data preview before sync trigger ,issue: datavane/tis#476
1 parent 8f92b84 commit cf5c04c

6 files changed

Lines changed: 12 additions & 12 deletions

File tree

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,7 @@ public PreviewRecords previewRecords(String tableName, QueryCriteria queryCriter
157157
Optional<Pair<String, List<String>>> transformer
158158
= transformerRules.map((trule) -> Pair.of(tableName, trule.relevantColKeys()));
159159

160-
this.startPipeline(readerCfg, transformer, (jobContainer) -> {
160+
this.startPipeline(tableName,readerCfg, transformer, (jobContainer) -> {
161161
jobContainer.setAttr(ThreadLocalRows.class, rows);
162162
});
163163

@@ -177,7 +177,7 @@ public PreviewRecords previewRecords(String tableName, QueryCriteria queryCriter
177177
* @param transformer
178178
* @param jobContainerSetter
179179
*/
180-
public void startPipeline(Configuration readerCfg, Optional<Pair<String, List<String>>> transformer, Consumer<JobContainer> jobContainerSetter) {
180+
public void startPipeline(String tableName,Configuration readerCfg, Optional<Pair<String, List<String>>> transformer, Consumer<JobContainer> jobContainerSetter) {
181181
Objects.requireNonNull(readerCfg);
182182
PerfTrace.getInstance(false, -1111, -1111, 0, false);
183183
Configuration allConf = IOUtils.loadResourceFromClasspath(DataxExecutor.class //
@@ -197,7 +197,7 @@ public void startPipeline(Configuration readerCfg, Optional<Pair<String, List<St
197197

198198

199199
Configuration c = Configuration.newDefault();
200-
c.set(CoreConstant.JOB_TRANSFORMER_NAME, transformer.map((p) -> p.getKey()).orElse("testTab"));
200+
c.set(CoreConstant.JOB_TRANSFORMER_NAME, tableName);
201201
transformer.ifPresent((tt) -> {
202202
Pair<String, List<String>> t = tt;
203203
c.set(CoreConstant.JOB_TRANSFORMER_RELEVANT_KEYS, t.getRight());

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,12 +51,12 @@ public WriterPluginMeta(String pluginKey, String streamwriterClass, Configuratio
5151
this.conf = conf;
5252
}
5353

54-
public static void realExecute(final String dataXName, final Configuration readerCfg
54+
public static void realExecute(final String dataXName,String tableName ,final Configuration readerCfg
5555
// , IDataXPluginMeta dataxReader
5656
, WriterPluginMeta writerPluginMeta
5757
, Optional<Pair<String, List<String>>> transformer, Optional<JarLoader> jarLoader) throws IllegalAccessException {
5858
realExecute(dataXName, writerPluginMeta, jarLoader)
59-
.startPipeline(readerCfg, transformer, (jobContainer) -> {
59+
.startPipeline(tableName,readerCfg, transformer, (jobContainer) -> {
6060
});
6161

6262
}

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

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,7 @@
2525
import com.alibaba.datax.common.spi.Writer;
2626
import com.alibaba.datax.common.util.Configuration;
2727
import com.google.common.collect.Lists;
28-
import org.apache.commons.lang3.Range;
2928

30-
import java.util.Collections;
3129
import java.util.List;
3230
import java.util.Objects;
3331

@@ -81,7 +79,8 @@ public void startWrite(RecordReceiver lineReceiver) {
8179
// gridRows.add(record);
8280
//recordToString(record);
8381
if (++readCount > queryCriteria.getPageSize()) {
84-
throw new IllegalStateException("readCount:" + readCount + " can not more than page size:" + queryCriteria.getPageSize());
82+
// throw new IllegalStateException("readCount:" + readCount + " can not more than page size:" + queryCriteria.getPageSize());
83+
return;
8584
}
8685
}
8786
}

tis-datax/tis-datax-local-executor/src/main/java/com/qlangtech/tis/plugin/datax/DataXPipelinePreviewProcessorExecutor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@
2727
import com.qlangtech.tis.datax.DataXJobSingleProcessorExecutor;
2828
import com.qlangtech.tis.datax.DataXJobSubmit.InstanceType;
2929
import com.qlangtech.tis.datax.DataXName;
30-
import com.qlangtech.tis.datax.DataxPrePostConsumer;
3130
import com.qlangtech.tis.datax.IDataXTaskRelevant;
3231
import com.qlangtech.tis.datax.TimeFormat;
3332
import com.qlangtech.tis.datax.preview.IPreviewRowsDataService;
@@ -42,6 +41,7 @@
4241
import com.qlangtech.tis.rpc.grpc.datax.preview.PreviewRowsDataCriteria.Builder;
4342
import com.qlangtech.tis.rpc.grpc.datax.preview.PreviewRowsDataResponse;
4443
import com.qlangtech.tis.rpc.grpc.datax.preview.StringValue;
44+
import com.qlangtech.tis.web.start.TisAppLaunch;
4545
import io.grpc.ConnectivityState;
4646
import io.grpc.ManagedChannel;
4747
import io.grpc.ManagedChannelBuilder;
@@ -260,7 +260,7 @@ protected String getMainClassName() {
260260

261261
@Override
262262
protected File getWorkingDirectory() {
263-
return DataXJobInfo.getDataXExecutorDir();
263+
return TisAppLaunch.isTestMock() ? new File("/opt/tis/tis-datax-executor") : DataXJobInfo.getDataXExecutorDir();
264264
}
265265

266266
@Override

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,7 @@ public static void realExecute(final String dataXName, final Configuration reade
159159
//" \"print\": true\n" +
160160
" \"" + Key.PATH + "\": \"" + writeFile.getParentFile().getAbsolutePath() + "\",\n" //
161161
+ " \"" + Key.FILE_NAME + "\": \"" + writeFile.getName() + "\"\n" + " }\n" + "}"));
162-
WriterPluginMeta.realExecute(dataXName, readerCfg, writerPluginMeta, transformer, Optional.empty());
162+
WriterPluginMeta.realExecute(dataXName, "testTabName", readerCfg, writerPluginMeta, transformer, Optional.empty());
163163

164164
// Objects.requireNonNull(readerCfg);
165165
// final JarLoader uberClassLoader = new JarLoader(new String[]{"."});

tis-incr/tis-flink-cdc-common/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFactory.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3434
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
3535
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
36+
import com.qlangtech.tis.trigger.util.UnCacheString;
3637
import com.qlangtech.tis.util.IPluginContext;
3738
import io.debezium.config.Field;
3839
import org.apache.commons.lang3.tuple.Pair;
@@ -113,7 +114,7 @@ public BasePGLikeDescriptor() {
113114
this.debeziumProps = Lists.newArrayList(
114115
Pair.of((opts) -> {
115116
opts.add(FIELD_KEY_SLOT_NAME, getSoltNameField()
116-
, new OverwriteProps().setDftVal(IPluginContext.getThreadLocalInstance().getCollectionName().getPipelineName()));
117+
, new OverwriteProps().setDftVal(new UnCacheString<>(() -> IPluginContext.getThreadLocalInstance().getCollectionName().getPipelineName())));
117118
}
118119
, (debeziumProperties, sourceFactory) -> {
119120
debeziumProperties.setProperty(getSoltNameField().name(), sourceFactory.slotName);

0 commit comments

Comments
 (0)