Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,16 @@
import org.apache.doris.common.AnalysisException;

import com.google.common.base.Strings;
import com.google.gson.annotations.SerializedName;

import java.io.StringWriter;

public class Separator {
private static final String HEX_STRING = "0123456789ABCDEF";

@SerializedName("os")
private final String oriSeparator;
@SerializedName("s")
private String separator;

public Separator(String separator) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,19 +26,29 @@
import org.apache.doris.load.loadv2.LoadTask;

import com.google.common.base.Strings;
import com.google.gson.annotations.SerializedName;

import java.util.List;

public class RoutineLoadDesc {
@SerializedName("cs")
private final Separator columnSeparator;
@SerializedName("ld")
private final Separator lineDelimiter;
@SerializedName("cols")
private final List<ImportColumnDesc> columnsInfo;
@SerializedName("pf")
private final Expr precedingFilter;
@SerializedName("f")
private final Expr filter;
@SerializedName("dc")
private final Expr deleteCondition;
@SerializedName("mt")
private LoadTask.MergeType mergeType;
// nullable
@SerializedName("pn")
private final PartitionNamesInfo partitionNamesInfo;
@SerializedName("sc")
private final String sequenceColName;

public RoutineLoadDesc(Separator columnSeparator, Separator lineDelimiter, List<ImportColumnDesc> columnsInfo,
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -943,7 +943,6 @@ public void alterRoutineLoadJob(AlterRoutineLoadCommand command) throws UserExce
+ command.getDataSourceProperties().getDataSourceType());
}
job.modifyProperties(command);
job.setRoutineLoadDesc(command.getRoutineLoadDesc());
}

public void replayAlterRoutineLoadJob(AlterRoutineLoadJobOperationLog log) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,6 @@
import org.apache.doris.rpc.RpcException;
import org.apache.doris.service.FrontendOptions;
import org.apache.doris.thrift.TFileCompressType;
import org.apache.doris.thrift.TPartialUpdateNewRowPolicy;
import org.apache.doris.transaction.TransactionState;
import org.apache.doris.transaction.TransactionStatus;

Expand All @@ -74,7 +73,6 @@
import com.google.gson.annotations.SerializedName;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.collections4.MapUtils;
import org.apache.commons.lang3.BooleanUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

Expand Down Expand Up @@ -778,6 +776,7 @@ public Map<String, String> getCustomProperties() {
@Override
public void modifyProperties(AlterRoutineLoadCommand command) throws UserException {
Map<String, String> jobProperties = command.getAnalyzedJobProperties();
validateCommonJobProperties(jobProperties);
KafkaDataSourceProperties dataSourceProperties = (KafkaDataSourceProperties) command.getDataSourceProperties();
if (null != dataSourceProperties) {
// if the partition offset is set by timestamp, convert it to real offset
Expand All @@ -791,9 +790,10 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti
}

modifyPropertiesInternal(jobProperties, dataSourceProperties);
setRoutineLoadDesc(command.getRoutineLoadDesc());

AlterRoutineLoadJobOperationLog log = new AlterRoutineLoadJobOperationLog(this.id,
jobProperties, dataSourceProperties);
jobProperties, dataSourceProperties, command.getRoutineLoadDesc());
Env.getCurrentEnv().getEditLog().logAlterRoutineLoadJob(log);
} finally {
writeUnlock();
Expand Down Expand Up @@ -881,17 +881,6 @@ private void modifyPropertiesInternal(Map<String, String> jobProperties,
Map<String, String> copiedJobProperties = Maps.newHashMap(jobProperties);
modifyCommonJobProperties(copiedJobProperties);
this.jobProperties.putAll(copiedJobProperties);
if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) {
this.isPartialUpdate = BooleanUtils.toBoolean(jobProperties.get(CreateRoutineLoadInfo.PARTIAL_COLUMNS));
}
if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) {
String policy = jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY);
if ("ERROR".equalsIgnoreCase(policy)) {
this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.ERROR;
} else {
this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND;
}
}
}
LOG.info("modify the properties of kafka routine load job: {}, jobProperties: {}, datasource properties: {}",
this.id, jobProperties, dataSourceProperties);
Expand Down Expand Up @@ -920,6 +909,7 @@ private void resetCloudProgress(Cloud.ResetRLProgressRequest.Builder builder) th
public void replayModifyProperties(AlterRoutineLoadJobOperationLog log) {
try {
modifyPropertiesInternal(log.getJobProperties(), (KafkaDataSourceProperties) log.getDataSourceProperties());
setRoutineLoadDesc(log.getRoutineLoadDesc());
} catch (UserException e) {
// should not happen
LOG.error("failed to replay modify kafka routine load job: {}", id, e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,6 @@
import org.apache.doris.persist.AlterRoutineLoadJobOperationLog;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.thrift.TFileCompressType;
import org.apache.doris.thrift.TPartialUpdateNewRowPolicy;
import org.apache.doris.transaction.TransactionState;
import org.apache.doris.transaction.TransactionStatus;

Expand All @@ -65,7 +64,6 @@
import com.google.gson.annotations.SerializedName;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.collections4.MapUtils;
import org.apache.commons.lang3.BooleanUtils;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

Expand Down Expand Up @@ -687,9 +685,10 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti
}

modifyPropertiesInternal(jobProperties, dataSourceProperties);
setRoutineLoadDesc(command.getRoutineLoadDesc());

AlterRoutineLoadJobOperationLog log = new AlterRoutineLoadJobOperationLog(this.id,
jobProperties, dataSourceProperties);
jobProperties, dataSourceProperties, command.getRoutineLoadDesc());
Env.getCurrentEnv().getEditLog().logAlterRoutineLoadJob(log);
} finally {
writeUnlock();
Expand All @@ -699,6 +698,7 @@ public void modifyProperties(AlterRoutineLoadCommand command) throws UserExcepti
private void modifyPropertiesInternal(Map<String, String> jobProperties,
KinesisDataSourceProperties dataSourceProperties)
throws UserException {
validateCommonJobProperties(jobProperties);
if (dataSourceProperties != null) {
List<Pair<String, String>> shardPositions = Lists.newArrayList();
Map<String, String> customKinesisProperties = Maps.newHashMap();
Expand Down Expand Up @@ -762,17 +762,6 @@ private void modifyPropertiesInternal(Map<String, String> jobProperties,
Map<String, String> copiedJobProperties = Maps.newHashMap(jobProperties);
modifyCommonJobProperties(copiedJobProperties);
this.jobProperties.putAll(copiedJobProperties);
if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_COLUMNS)) {
this.isPartialUpdate = BooleanUtils.toBoolean(jobProperties.get(CreateRoutineLoadInfo.PARTIAL_COLUMNS));
}
if (jobProperties.containsKey(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY)) {
String policy = jobProperties.get(CreateRoutineLoadInfo.PARTIAL_UPDATE_NEW_KEY_POLICY);
if ("ERROR".equalsIgnoreCase(policy)) {
this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.ERROR;
} else {
this.partialUpdateNewKeyPolicy = TPartialUpdateNewRowPolicy.APPEND;
}
}
}
LOG.info("modify the properties of kinesis routine load job: {}, jobProperties: {}, datasource properties: {}",
this.id, jobProperties, dataSourceProperties);
Expand All @@ -783,6 +772,7 @@ public void replayModifyProperties(AlterRoutineLoadJobOperationLog log) {
try {
modifyPropertiesInternal(log.getJobProperties(),
(KinesisDataSourceProperties) log.getDataSourceProperties());
setRoutineLoadDesc(log.getRoutineLoadDesc());
} catch (UserException e) {
LOG.error("failed to replay modify kinesis routine load job: {}", id, e);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import org.apache.doris.common.io.Text;
import org.apache.doris.common.io.Writable;
import org.apache.doris.load.RoutineLoadDesc;
import org.apache.doris.load.routineload.AbstractDataSourceProperties;
import org.apache.doris.persist.gson.GsonUtils;

Expand All @@ -37,12 +38,20 @@ public class AlterRoutineLoadJobOperationLog implements Writable {
private Map<String, String> jobProperties;
@SerializedName(value = "dataSourceProperties")
private AbstractDataSourceProperties dataSourceProperties;
@SerializedName(value = "routineLoadDesc")
private RoutineLoadDesc routineLoadDesc;

public AlterRoutineLoadJobOperationLog(long jobId, Map<String, String> jobProperties,
AbstractDataSourceProperties dataSourceProperties) {
this(jobId, jobProperties, dataSourceProperties, null);
}

public AlterRoutineLoadJobOperationLog(long jobId, Map<String, String> jobProperties,
AbstractDataSourceProperties dataSourceProperties, RoutineLoadDesc routineLoadDesc) {
this.jobId = jobId;
this.jobProperties = jobProperties;
this.dataSourceProperties = dataSourceProperties;
this.routineLoadDesc = routineLoadDesc;
}

public long getJobId() {
Expand All @@ -57,6 +66,10 @@ public AbstractDataSourceProperties getDataSourceProperties() {
return dataSourceProperties;
}

public RoutineLoadDesc getRoutineLoadDesc() {
return routineLoadDesc;
}

public static AlterRoutineLoadJobOperationLog read(DataInput in) throws IOException {
String json = Text.readString(in);
return GsonUtils.GSON.fromJson(json, AlterRoutineLoadJobOperationLog.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.datasource.InternalCatalog;
import org.apache.doris.datasource.kafka.KafkaUtil;
import org.apache.doris.datasource.property.fileformat.CsvFileFormatProperties;
import org.apache.doris.load.RoutineLoadDesc;
import org.apache.doris.load.loadv2.LoadTask;
import org.apache.doris.load.routineload.kafka.KafkaConfiguration;
Expand All @@ -40,10 +41,13 @@
import org.apache.doris.load.routineload.kafka.KafkaRoutineLoadJob;
import org.apache.doris.load.routineload.kafka.KafkaTaskInfo;
import org.apache.doris.mysql.privilege.MockedAuth;
import org.apache.doris.nereids.trees.plans.commands.AlterRoutineLoadCommand;
import org.apache.doris.nereids.trees.plans.commands.info.CreateRoutineLoadInfo;
import org.apache.doris.nereids.trees.plans.commands.info.LabelNameInfo;
import org.apache.doris.nereids.trees.plans.commands.load.LoadProperty;
import org.apache.doris.nereids.trees.plans.commands.load.LoadSeparator;
import org.apache.doris.persist.AlterRoutineLoadJobOperationLog;
import org.apache.doris.persist.EditLog;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.thrift.TResourceInfo;
import org.apache.doris.thrift.TRoutineLoadTask;
Expand All @@ -58,9 +62,14 @@
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.MockedStatic;
import org.mockito.Mockito;

import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.DataInputStream;
import java.io.DataOutputStream;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
Expand Down Expand Up @@ -272,6 +281,84 @@ public void testUpdateProgressWarnsWhenReadCommittedTaskHasZeroRowsAndLag() thro
Assert.assertTrue(otherMsg.contains("some records may be in uncommitted transactions"));
}

@Test
public void testAlterPersistsLoadDescAndCsvPropertiesForReplay() throws Exception {
KafkaRoutineLoadJob leader = createPausedJob();
KafkaRoutineLoadJob follower = createPausedJob();
RoutineLoadDesc originalDesc = new RoutineLoadDesc(new Separator("|", "|"), null, null,
null, null, null, null, LoadTask.MergeType.APPEND, "original_sequence");
leader.setRoutineLoadDesc(originalDesc);
follower.setRoutineLoadDesc(originalDesc);

Map<String, String> jobProperties = Maps.newHashMap();
jobProperties.put(CsvFileFormatProperties.PROP_ENCLOSE, "\"");
jobProperties.put(CsvFileFormatProperties.PROP_ESCAPE, "\\");
jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, "true");
RoutineLoadDesc delta = new RoutineLoadDesc(null, new Separator("\n", "\\n"), null,
null, null, null, null, LoadTask.MergeType.APPEND, null);
AlterRoutineLoadCommand command = Mockito.mock(AlterRoutineLoadCommand.class);
Mockito.when(command.getAnalyzedJobProperties()).thenReturn(jobProperties);
Mockito.when(command.getDataSourceProperties()).thenReturn(null);
Mockito.when(command.getRoutineLoadDesc()).thenReturn(delta);

Env env = Mockito.mock(Env.class);
EditLog editLog = Mockito.mock(EditLog.class);
AlterRoutineLoadJobOperationLog alterLog;
try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class)) {
envStatic.when(Env::getCurrentEnv).thenReturn(env);
Mockito.when(env.getEditLog()).thenReturn(editLog);

leader.modifyProperties(command);

ArgumentCaptor<AlterRoutineLoadJobOperationLog> logCaptor =
ArgumentCaptor.forClass(AlterRoutineLoadJobOperationLog.class);
Mockito.verify(editLog).logAlterRoutineLoadJob(logCaptor.capture());
alterLog = logCaptor.getValue();
}

Assert.assertSame(delta, alterLog.getRoutineLoadDesc());
Assert.assertEquals(jobProperties, alterLog.getJobProperties());
assertAlterState(leader);

follower.replayModifyProperties(alterLog);
assertAlterState(follower);

assertAlterState(imageRoundTrip(leader));
assertAlterState(imageRoundTrip(follower));
}

private static KafkaRoutineLoadJob createPausedJob() {
KafkaRoutineLoadJob job = new KafkaRoutineLoadJob(1L, "job1", 1L,
1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN);
Deencapsulation.setField(job, "state", RoutineLoadJob.JobState.PAUSED);
return job;
}

private static void assertAlterState(RoutineLoadJob job) {
Assert.assertEquals("|", job.getColumnSeparator().getSeparator());
Assert.assertEquals("\n", job.getLineDelimiter().getSeparator());
Assert.assertEquals("original_sequence", job.getSequenceCol());
Assert.assertEquals((byte) '"', job.getEnclose());
Assert.assertEquals((byte) '\\', job.getEscape());
Assert.assertTrue(job.getEmptyFieldAsNull());
Assert.assertEquals(Boolean.TRUE, Deencapsulation.getField(job, "emptyFieldAsNull"));

Map<String, String> persistedJobProperties = Deencapsulation.getField(job, "jobProperties");
Assert.assertEquals("\"", persistedJobProperties.get(CsvFileFormatProperties.PROP_ENCLOSE));
Assert.assertEquals("\\", persistedJobProperties.get(CsvFileFormatProperties.PROP_ESCAPE));
Assert.assertEquals("true", persistedJobProperties.get(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL));
}

private static RoutineLoadJob imageRoundTrip(RoutineLoadJob routineLoadJob) throws Exception {
ByteArrayOutputStream bytes = new ByteArrayOutputStream();
try (DataOutputStream out = new DataOutputStream(bytes)) {
routineLoadJob.write(out);
}
try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes.toByteArray()))) {
return RoutineLoadJob.read(in);
}
}

@Test
public void testDisplayCustomPropertiesMasksKafkaSecrets() {
KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L,
Expand Down
Loading
Loading