github-actions[bot] commented on code in PR #67512:
URL: https://github.com/apache/doris/pull/67512#discussion_r4227992808


##########
fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java:
##########
@@ -272,6 +272,41 @@ public void 
testUpdateProgressWarnsWhenReadCommittedTaskHasZeroRowsAndLag() thro
         Assert.assertTrue(otherMsg.contains("some records may be in 
uncommitted transactions"));
     }
 
+    @Test
+    public void testMultiTableCsvTaskIncludesParserProperties() throws 
Exception {
+        RoutineLoadManager routineLoadManager = 
Mockito.mock(RoutineLoadManager.class);
+        Env env = Mockito.mock(Env.class);
+        InternalCatalog internalCatalog = Mockito.mock(InternalCatalog.class);
+        Database database = Mockito.mock(Database.class);
+
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class)) {
+            envStatic.when(Env::getCurrentEnv).thenReturn(env);
+            
envStatic.when(Env::getCurrentInternalCatalog).thenReturn(internalCatalog);
+            
Mockito.when(env.getRoutineLoadManager()).thenReturn(routineLoadManager);
+            
Mockito.when(internalCatalog.getDbOrMetaException(1L)).thenReturn(database);
+            Mockito.when(database.getFullName()).thenReturn("db1");
+
+            KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, 
"multi_table_job", 1L,
+                    "127.0.0.1:9020", "topic1", UserIdentity.ADMIN, true);
+            Deencapsulation.setField(routineLoadJob, "enclose", (byte) '^');
+            Deencapsulation.setField(routineLoadJob, "escape", (byte) '?');
+            Map<String, String> jobProperties = 
Deencapsulation.getField(routineLoadJob, "jobProperties");
+            
jobProperties.put(CsvFileFormatProperties.PROP_EMPTY_FIELD_AS_NULL, "true");

Review Comment:
   [P1] Import `CsvFileFormatProperties` in this test. The test is in 
`org.apache.doris.load.routineload`, while this class is in 
`org.apache.doris.datasource.property.fileformat`, and the new reference has no 
import. FE test compilation fails with an unresolved symbol until the import is 
added.



##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaTaskInfo.java:
##########
@@ -121,6 +121,11 @@ public TRoutineLoadTask createRoutineLoadTask() throws 
UserException {
             tRoutineLoadTask.setFormat(TFileFormatType.FORMAT_JSON);
         } else {
             tRoutineLoadTask.setFormat(TFileFormatType.FORMAT_CSV_PLAIN);
+            if (isMultiTable) {
+                tRoutineLoadTask.setEnclose(routineLoadJob.getEnclose());

Review Comment:
   [P1] Use current CSV quote settings when creating multi-table tasks. 
`getEnclose()` and `getEscape()` read fields that Gson omits from Routine Load 
job persistence, and `ALTER ROUTINE LOAD` updates only `jobProperties`; after 
an FE restart these setters send zero, and after ALTER they send the previous 
bytes. BE forwards those bytes into every per-table CSV scan, so quoted or 
escaped rows can be split or parsed incorrectly. Derive the values from the 
persisted/current properties, or synchronize the fields on both replay and 
ALTER, and cover those transitions in a test.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to