This is an automated email from the ASF dual-hosted git repository.
JNSimba pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 193342aaf3f [fix](streamingjob) Support precreated streaming job
targets (#66950)
193342aaf3f is described below
commit 193342aaf3fc2121c726ac07623c752c01cce1e7
Author: wudi <[email protected]>
AuthorDate: Fri Aug 28 16:40:43 2026 +0800
[fix](streamingjob) Support precreated streaming job targets (#66950)
### What problem does this PR solve?
Issue Number: None
Related PR: None
Problem Summary: Streaming Job converted every source JDBC column and
generated target DDL before checking whether the mapped Doris target
table already existed. An unsupported source type therefore prevented
CREATE JOB even when the user had already created a compatible target
table. This change resolves the mapped target table first and skips
source-to-Doris type conversion only when that target already exists.
Source table identity, primary-key validation, excluded-primary-key
validation, and automatic table creation for missing targets remain
unchanged.
---
.../insert/streaming/StreamingInsertJob.java | 17 ++-
.../apache/doris/job/util/StreamingJobUtils.java | 40 ++++--
.../doris/job/util/StreamingJobUtilsTest.java | 46 +++++++
...st_streaming_postgres_job_precreated_target.out | 4 +
...streaming_postgres_job_precreated_target.groovy | 134 +++++++++++++++++++++
5 files changed, 218 insertions(+), 23 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
index 6b9c2234012..185aa720c4a 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java
@@ -18,7 +18,6 @@
package org.apache.doris.job.extensions.insert.streaming;
import org.apache.doris.analysis.UserIdentity;
-import org.apache.doris.catalog.Database;
import org.apache.doris.catalog.Env;
import org.apache.doris.cloud.catalog.CloudEnv;
import org.apache.doris.cloud.proto.Cloud;
@@ -103,6 +102,7 @@ import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
import java.util.TreeMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicLong;
@@ -301,18 +301,15 @@ public class StreamingInsertJob extends
AbstractJob<StreamingJobSchedulerTask, M
private List<String> createTableIfNotExists() throws Exception {
List<String> syncTbls = new ArrayList<>();
Map<String, String> effectiveSourceProperties =
buildConvertedSourceProperties(sourceProperties);
- // Key: source table name; Value: CreateTableCommand for the Doris
target table.
- // The two names differ when "table.<src>.target_table" is configured.
- LinkedHashMap<String, CreateTableCommand> createTblCmds =
+ // Key: source table name; Value: CREATE TABLE command, or empty if
the target already exists.
+ // The source and target table names differ when
"table.<src>.target_table" is configured.
+ LinkedHashMap<String, Optional<CreateTableCommand>> createTblCmds =
StreamingJobUtils.generateCreateTableCmds(targetDb,
dataSourceType, effectiveSourceProperties,
targetProperties);
- Database db =
Env.getCurrentEnv().getInternalCatalog().getDbNullable(targetDb);
- Preconditions.checkNotNull(db, "target database %s does not exist",
targetDb);
- for (Map.Entry<String, CreateTableCommand> entry :
createTblCmds.entrySet()) {
+ for (Map.Entry<String, Optional<CreateTableCommand>> entry :
createTblCmds.entrySet()) {
String srcTable = entry.getKey();
- CreateTableCommand createTblCmd = entry.getValue();
- if
(!db.isTableExist(createTblCmd.getCreateTableInfo().getTableName())) {
- createTblCmd.run(ConnectContext.get(), null);
+ if (entry.getValue().isPresent()) {
+ entry.getValue().get().run(ConnectContext.get(), null);
}
// Use the upstream table name so CDC monitors the correct source
table.
syncTbls.add(srcTable);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
b/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
index 0877e43a698..c36b08b2edf 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/job/util/StreamingJobUtils.java
@@ -25,6 +25,7 @@ import org.apache.doris.catalog.KeysType;
import org.apache.doris.catalog.PrimitiveType;
import org.apache.doris.catalog.ScalarType;
import org.apache.doris.catalog.Table;
+import org.apache.doris.catalog.info.TableNameInfo;
import org.apache.doris.cloud.system.CloudSystemInfoService;
import org.apache.doris.common.Config;
import org.apache.doris.common.DdlException;
@@ -358,15 +359,15 @@ public class StreamingJobUtils {
*
* <p>Returns a {@link LinkedHashMap} whose key is the <b>source</b>
(upstream) table name and
* whose value is the corresponding {@link CreateTableCommand} that
creates the Doris target
- * table (which may have a different name when {@code
table.<src>.target_table} is configured).
- * Callers must use the map key as the upstream source table identifier
for CDC monitoring and
- * the {@link CreateTableCommand} value for the actual DDL execution.
+ * table (which may have a different name when {@code
table.<src>.target_table} is configured),
+ * or empty when the target table already exists. Callers must use the map
key as the upstream
+ * source table identifier for CDC monitoring.
*/
- public static LinkedHashMap<String, CreateTableCommand>
generateCreateTableCmds(String targetDb,
+ public static LinkedHashMap<String, Optional<CreateTableCommand>>
generateCreateTableCmds(String targetDb,
DataSourceType sourceType,
Map<String, String> properties, Map<String, String>
targetProperties)
throws JobException {
- LinkedHashMap<String, CreateTableCommand> createtblCmds = new
LinkedHashMap<>();
+ LinkedHashMap<String, Optional<CreateTableCommand>> createtblCmds =
new LinkedHashMap<>();
String includeTables =
properties.get(DataSourceConfigKeys.INCLUDE_TABLES);
String excludeTables =
properties.get(DataSourceConfigKeys.EXCLUDE_TABLES);
List<String> includeTablesList = new ArrayList<>();
@@ -385,6 +386,8 @@ public class StreamingJobUtils {
if (tablesNameList.isEmpty()) {
throw new JobException("No tables found in database " +
database);
}
+ Database targetDatabase =
Env.getCurrentEnv().getInternalCatalog().getDbNullable(targetDb);
+ Preconditions.checkNotNull(targetDatabase, "target database %s
does not exist", targetDb);
Map<String, String> tableCreateProperties =
getTableCreateProperties(targetProperties);
List<String> noPrimaryKeyTables = new ArrayList<>();
@@ -404,21 +407,33 @@ public class StreamingJobUtils {
}
List<String> primaryKeys = jdbcClient.getPrimaryKeys(database,
table);
- List<Column> columns = getColumns(jdbcClient, database, table,
primaryKeys);
if (primaryKeys.isEmpty()) {
noPrimaryKeyTables.add(table);
}
// Resolve target (Doris) table name; defaults to source table
name if not configured
- String targetTableName = properties.getOrDefault(
+ String targetTableName = new TableNameInfo(targetDb,
properties.getOrDefault(
DataSourceConfigKeys.TABLE + "." + table + "."
+
DataSourceConfigKeys.TABLE_TARGET_TABLE_SUFFIX,
- table).trim();
+ table).trim()).getTbl();
// Validate and apply exclude_columns for this table
Set<String> excludeColumns = parseExcludeColumns(properties,
table);
+ if (targetDatabase.isTableExist(targetTableName)) {
+ if (!excludeColumns.isEmpty()) {
+ Set<String> columnNames =
jdbcClient.getJdbcColumnsInfo(database, table).stream()
+ .map(field -> field.getColumnName())
+ .collect(Collectors.toSet());
+ validateExcludeColumns(excludeColumns, table,
columnNames, primaryKeys);
+ }
+ createtblCmds.put(table, Optional.empty());
+ continue;
+ }
+
+ List<Column> columns = getColumns(jdbcClient, database, table,
primaryKeys);
if (!excludeColumns.isEmpty()) {
- validateExcludeColumns(excludeColumns, table, columns,
primaryKeys);
+ Set<String> columnNames =
columns.stream().map(Column::getName).collect(Collectors.toSet());
+ validateExcludeColumns(excludeColumns, table, columnNames,
primaryKeys);
columns = columns.stream()
.filter(col ->
!excludeColumns.contains(col.getName()))
.collect(Collectors.toList());
@@ -461,7 +476,7 @@ public class StreamingJobUtils {
);
CreateTableCommand createtblCmd = new
CreateTableCommand(Optional.empty(), createtblInfo);
// Key: source (PG/MySQL) table name; Value: command that
creates the Doris target table
- createtblCmds.put(table, createtblCmd);
+ createtblCmds.put(table, Optional.of(createtblCmd));
}
if (createtblCmds.isEmpty()) {
throw new JobException("Can not found match table in database
" + database);
@@ -623,10 +638,9 @@ public class StreamingJobUtils {
}
private static void validateExcludeColumns(Set<String> excludeColumns,
String tableName,
- List<Column> columns, List<String> primaryKeys) throws
JobException {
- Set<String> colNames =
columns.stream().map(Column::getName).collect(Collectors.toSet());
+ Set<String> columnNames, List<String> primaryKeys) throws
JobException {
for (String col : excludeColumns) {
- if (!colNames.contains(col)) {
+ if (!columnNames.contains(col)) {
throw new JobException(String.format(
"exclude_columns validation failed: column '%s' does
not exist in table '%s'",
col, tableName));
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/job/util/StreamingJobUtilsTest.java
b/fe/fe-core/src/test/java/org/apache/doris/job/util/StreamingJobUtilsTest.java
index b2e2b0a56be..e8279cf3d1d 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/job/util/StreamingJobUtilsTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/job/util/StreamingJobUtilsTest.java
@@ -18,12 +18,18 @@
package org.apache.doris.job.util;
import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Database;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.KeysType;
+import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.PrimitiveType;
import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.datasource.InternalCatalog;
import org.apache.doris.datasource.jdbc.client.JdbcClient;
import org.apache.doris.job.cdc.DataSourceConfigKeys;
import org.apache.doris.job.common.DataSourceType;
import org.apache.doris.job.exception.JobException;
+import org.apache.doris.qe.GlobalVariable;
import org.junit.Assert;
import org.junit.Before;
@@ -270,4 +276,44 @@ public class StreamingJobUtilsTest {
Mockito.verify(jdbcClient).closeClient();
}
}
+
+ @Test
+ public void
testGenerateCreateTableCmdsFindsMixedCasePrecreatedTargetWhenStoredLowerCase()
throws Exception {
+ Map<String, String> properties = new HashMap<>();
+ properties.put(DataSourceConfigKeys.SCHEMA, "source_db");
+ properties.put(DataSourceConfigKeys.TABLE + ".source_table."
+ + DataSourceConfigKeys.TABLE_TARGET_TABLE_SUFFIX,
"MixedTarget");
+
+ Database targetDatabase = new Database(1L, "target_db");
+ targetDatabase.registerTable(new OlapTable(2L, "mixedtarget", new
ArrayList<>(), KeysType.UNIQUE_KEYS,
+ null, null));
+ Env env = Mockito.mock(Env.class);
+ InternalCatalog internalCatalog = Mockito.mock(InternalCatalog.class);
+ Mockito.when(env.getInternalCatalog()).thenReturn(internalCatalog);
+
Mockito.when(internalCatalog.getDbNullable("target_db")).thenReturn(targetDatabase);
+ Mockito.when(jdbcClient.getTablesNameList("source_db"))
+ .thenReturn(Arrays.asList("source_table"));
+ Mockito.when(jdbcClient.getPrimaryKeys("source_db", "source_table"))
+ .thenReturn(Arrays.asList("id"));
+ Mockito.when(jdbcClient.getColumnsFromJdbc("source_db",
"source_table"))
+ .thenReturn(Arrays.asList(
+ new Column("id",
ScalarType.createType(PrimitiveType.INT)),
+ new Column("unsupported_col", new
ScalarType(PrimitiveType.UNSUPPORTED))));
+
+ int originalLowerCaseTableNames = GlobalVariable.lowerCaseTableNames;
+ GlobalVariable.lowerCaseTableNames = 1;
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class);
+ MockedStatic<StreamingJobUtils> utils =
Mockito.mockStatic(StreamingJobUtils.class,
+ Mockito.CALLS_REAL_METHODS)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ utils.when(() ->
StreamingJobUtils.getJdbcClient(DataSourceType.POSTGRES, properties))
+ .thenReturn(jdbcClient);
+
+ Assert.assertFalse(StreamingJobUtils.generateCreateTableCmds(
+ "target_db", DataSourceType.POSTGRES, properties, new
HashMap<>())
+ .get("source_table").isPresent());
+ } finally {
+ GlobalVariable.lowerCaseTableNames = originalLowerCaseTableNames;
+ }
+ }
}
diff --git
a/regression-test/data/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.out
b/regression-test/data/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.out
new file mode 100644
index 00000000000..45a44ef6913
--- /dev/null
+++
b/regression-test/data/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.out
@@ -0,0 +1,4 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !select_precreated_target --
+1 ready
+
diff --git
a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.groovy
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.groovy
new file mode 100644
index 00000000000..9e64974fc58
--- /dev/null
+++
b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_postgres_job_precreated_target.groovy
@@ -0,0 +1,134 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.awaitility.Awaitility
+
+import static java.util.concurrent.TimeUnit.SECONDS
+
+suite("test_streaming_postgres_job_precreated_target",
+ "p0,external,pg,external_docker,external_docker_pg,nondatalake") {
+ def jobName = "test_streaming_postgres_job_precreated_target"
+ def currentDb = (sql "select database()")[0][0]
+ def pgDB = "postgres"
+ def pgSchema = "cdc_test"
+ def pgUser = "postgres"
+ def pgPassword = "123456"
+
+ sql """DROP JOB IF EXISTS where jobname = '${jobName}'"""
+ sql """DROP JOB IF EXISTS where jobname = '${jobName}_invalid_exclude'"""
+ sql """drop table if exists ${currentDb}.streaming_precreated_enum_target
force"""
+
+ String enabled = context.config.otherConfigs.get("enableJdbcTest")
+ if (enabled != null && enabled.equalsIgnoreCase("true")) {
+ String pgPort = context.config.otherConfigs.get("pg_14_port")
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+ String s3Endpoint = getS3Endpoint()
+ String bucket = getS3BucketName()
+ String driverUrl =
"https://${bucket}.${s3Endpoint}/regression/jdbc_driver/postgresql-42.5.0.jar"
+
+ connect("${pgUser}", "${pgPassword}",
"jdbc:postgresql://${externalEnvIp}:${pgPort}/${pgDB}") {
+ sql """DROP TABLE IF EXISTS
${pgDB}.${pgSchema}.streaming_precreated_enum_source"""
+ sql """DROP TYPE IF EXISTS
${pgSchema}.streaming_precreated_target_enum"""
+ sql """CREATE TYPE ${pgSchema}.streaming_precreated_target_enum AS
ENUM ('ready', 'done')"""
+ sql """
+ CREATE TABLE
${pgDB}.${pgSchema}.streaming_precreated_enum_source (
+ id INTEGER PRIMARY KEY,
+ searchable ${pgSchema}.streaming_precreated_target_enum
+ )
+ """
+ sql """
+ INSERT INTO
${pgDB}.${pgSchema}.streaming_precreated_enum_source
+ VALUES (1, 'ready')
+ """
+ }
+
+ sql """
+ CREATE TABLE ${currentDb}.streaming_precreated_enum_target (
+ id INT NOT NULL,
+ searchable STRING NULL
+ ) ENGINE=OLAP
+ UNIQUE KEY(id)
+ DISTRIBUTED BY HASH(id) BUCKETS AUTO
+ PROPERTIES ("replication_num" = "1")
+ """
+
+ sql """
+ CREATE JOB ${jobName}
+ ON STREAMING
+ FROM POSTGRES (
+ "jdbc_url" =
"jdbc:postgresql://${externalEnvIp}:${pgPort}/${pgDB}",
+ "driver_url" = "${driverUrl}",
+ "driver_class" = "org.postgresql.Driver",
+ "user" = "${pgUser}",
+ "password" = "${pgPassword}",
+ "database" = "${pgDB}",
+ "schema" = "${pgSchema}",
+ "include_tables" = "streaming_precreated_enum_source",
+ "offset" = "initial",
+ "table.streaming_precreated_enum_source.target_table" =
"streaming_precreated_enum_target"
+ )
+ TO DATABASE ${currentDb} (
+ "table.create.properties.replication_num" = "1"
+ )
+ """
+
+ try {
+ Awaitility.await().atMost(300, SECONDS).pollInterval(2,
SECONDS).until {
+ def rows = sql """
+ SELECT id, searchable
+ FROM ${currentDb}.streaming_precreated_enum_target
+ WHERE id = 1
+ """
+ rows.size() == 1 && rows[0][0].toString() == "1" && rows[0][1]
== "ready"
+ }
+ } catch (Exception ex) {
+ log.info("show job: " + (sql """select * from
jobs("type"="insert") where Name='${jobName}'"""))
+ log.info("show task: " + (sql """select * from
tasks("type"="insert") where JobName='${jobName}'"""))
+ throw ex
+ }
+
+ order_qt_select_precreated_target """
+ SELECT id, searchable FROM
${currentDb}.streaming_precreated_enum_target
+ """
+
+ sql """DROP JOB IF EXISTS where jobname = '${jobName}'"""
+
+ test {
+ sql """
+ CREATE JOB ${jobName}_invalid_exclude
+ ON STREAMING
+ FROM POSTGRES (
+ "jdbc_url" =
"jdbc:postgresql://${externalEnvIp}:${pgPort}/${pgDB}",
+ "driver_url" = "${driverUrl}",
+ "driver_class" = "org.postgresql.Driver",
+ "user" = "${pgUser}",
+ "password" = "${pgPassword}",
+ "database" = "${pgDB}",
+ "schema" = "${pgSchema}",
+ "include_tables" = "streaming_precreated_enum_source",
+ "offset" = "initial",
+ "table.streaming_precreated_enum_source.target_table" =
"streaming_precreated_enum_target",
+ "table.streaming_precreated_enum_source.exclude_columns" =
"missing_searchable"
+ )
+ TO DATABASE ${currentDb} (
+ "table.create.properties.replication_num" = "1"
+ )
+ """
+ exception "exclude_columns validation failed: column
'missing_searchable' does not exist"
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]