This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12526-1ea469f5cc8a2ac5ccce8e7f90b4dc7f3b8bf18a in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 483aa86c096a2d72eba28ac85eb8d3a5c71e9295 Author: Jast <[email protected]> AuthorDate: Fri Oct 2 08:04:37 2026 +0000 [Fix][Connector-V2] Support special characters in column names for ClickHouse sink named-parameter SQL (#12526) Co-authored-by: jast <[email protected]> --- .../executor/FieldNamedPreparedStatement.java | 82 ++++++++++++++-- .../sink/client/executor/SqlUtilsTest.java | 103 +++++++++++++++++++++ .../seatunnel/clickhouse/ClickhouseIT.java | 79 ++++++++++++++++ .../clickhouse_special_columns_to_clickhouse.conf | 50 ++++++++++ 4 files changed, 307 insertions(+), 7 deletions(-) diff --git a/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/FieldNamedPreparedStatement.java b/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/FieldNamedPreparedStatement.java index da543ad5db..6b0e03fdef 100644 --- a/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/FieldNamedPreparedStatement.java +++ b/seatunnel-connectors-v2/connector-clickhouse/src/main/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/FieldNamedPreparedStatement.java @@ -642,7 +642,7 @@ public class FieldNamedPreparedStatement implements PreparedStatement { } } else { HashMap<String, List<Integer>> parameterMap = new HashMap<>(); - parsedSQL = parseNamedStatement(sql, parameterMap, parameterExpression); + parsedSQL = parseNamedStatement(sql, parameterMap, parameterExpression, fieldNames); // currently, the statements must contain all the field parameters checkArgument(parameterMap.size() >= fieldNames.length); for (int i = 0; i < fieldNames.length; i++) { @@ -658,13 +658,28 @@ public class FieldNamedPreparedStatement implements PreparedStatement { } public static String parseNamedStatement(String sql, Map<String, List<Integer>> paramMap) { - return parseNamedStatement(sql, paramMap, field -> "?"); + return parseNamedStatement(sql, paramMap, field -> "?", null); + } + + /** + * Parses named parameters (":name") in the given statement. + * + * <p>ClickHouse SQL templates use sink field names as named parameters. This overload uses that + * field list as the allow-list for names containing characters outside the legacy ClickHouse + * identifier scan, such as '#', spaces, or '-'. It intentionally keeps the old scan for normal + * identifiers (including dotted ClickHouse identifiers), still rejects empty names, and does + * not attempt to parse quoted SQL text or placeholders containing '?'. + */ + public static String parseNamedStatement( + String sql, Map<String, List<Integer>> paramMap, String[] knownParameterNames) { + return parseNamedStatement(sql, paramMap, field -> "?", knownParameterNames); } private static String parseNamedStatement( String sql, Map<String, List<Integer>> paramMap, - Function<String, String> parameterExpression) { + Function<String, String> parameterExpression, + String[] knownParameterNames) { StringBuilder parsedSql = new StringBuilder(); int fieldIndex = 1; // SQL statement parameter index starts from 1 int length = sql.length(); @@ -672,10 +687,15 @@ public class FieldNamedPreparedStatement implements PreparedStatement { char c = sql.charAt(i); if (':' == c) { int j = i + 1; - while (j < length - && (Character.isJavaIdentifierPart(sql.charAt(j)) - || ".".equals(String.valueOf(sql.charAt(j))))) { - j++; + String knownParameter = matchKnownParameter(sql, i + 1, knownParameterNames); + if (knownParameter != null) { + j = i + 1 + knownParameter.length(); + } else { + while (j < length + && (Character.isJavaIdentifierPart(sql.charAt(j)) + || ".".equals(String.valueOf(sql.charAt(j))))) { + j++; + } } String parameterName = sql.substring(i + 1, j); checkArgument( @@ -691,4 +711,52 @@ public class FieldNamedPreparedStatement implements PreparedStatement { } return parsedSql.toString(); } + + /** + * Returns the longest known parameter name that starts exactly at {@code offset} and is not + * immediately followed by more legacy identifier characters, or {@code null} when there is no + * such match (the caller then falls back to the identifier scan). This must stay stricter than + * a general SQL tokenizer; it only recognizes names from the current statement's field list. + */ + private static String matchKnownParameter( + String sql, int offset, String[] knownParameterNames) { + if (knownParameterNames == null || knownParameterNames.length == 0) { + return null; + } + if (offset >= sql.length() || !Character.isJavaIdentifierPart(sql.charAt(offset))) { + return null; + } + String best = null; + for (String name : knownParameterNames) { + if (name == null + || name.isEmpty() + || name.indexOf(':') >= 0 + || isLegacyIdentifier(name) + || !sql.startsWith(name, offset)) { + continue; + } + int end = offset + name.length(); + if (end < sql.length() + && (Character.isJavaIdentifierPart(sql.charAt(end)) + || ".".equals(String.valueOf(sql.charAt(end))))) { + // ClickHouse legacy parsing treats dots as part of a token, so keep prefix matches + // from stealing the front of a longer parameter such as ":db.name_extra". + continue; + } + if (best == null || name.length() > best.length()) { + best = name; + } + } + return best; + } + + private static boolean isLegacyIdentifier(String name) { + for (int i = 0; i < name.length(); i++) { + char c = name.charAt(i); + if (!Character.isJavaIdentifierPart(c) && c != '.') { + return false; + } + } + return true; + } } diff --git a/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/SqlUtilsTest.java b/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/SqlUtilsTest.java new file mode 100644 index 0000000000..ea1ee65462 --- /dev/null +++ b/seatunnel-connectors-v2/connector-clickhouse/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/sink/client/executor/SqlUtilsTest.java @@ -0,0 +1,103 @@ +/* + * 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. + */ + +package org.apache.seatunnel.connectors.seatunnel.clickhouse.sink.client.executor; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** Tests generated ClickHouse sink SQL for special characters in column names. */ +class SqlUtilsTest { + + @Test + void parseNamedStatementWithSpecialClickHouseColumnNames() { + String insertSQL = + SqlUtils.getInsertIntoStatement( + "sink_table", + new String[] {"GLREG", "GLREG#", "MY COL", "COL-1", "db.name"}); + String[] fieldNames = {"GLREG", "GLREG#", "MY COL", "COL-1", "db.name"}; + + Map<String, List<Integer>> parameterMap = new HashMap<>(); + String parsedSQL = + FieldNamedPreparedStatement.parseNamedStatement( + insertSQL, parameterMap, fieldNames); + + Assertions.assertEquals(5, parameterMap.size()); + Assertions.assertEquals("[1]", parameterMap.get("GLREG").toString()); + Assertions.assertEquals("[2]", parameterMap.get("GLREG#").toString()); + Assertions.assertEquals("[3]", parameterMap.get("MY COL").toString()); + Assertions.assertEquals("[4]", parameterMap.get("COL-1").toString()); + Assertions.assertEquals("[5]", parameterMap.get("db.name").toString()); + Assertions.assertEquals( + "INSERT INTO sink_table (\"GLREG\", \"GLREG#\", \"MY COL\", \"COL-1\", \"db.name\") VALUES (?, ?, ?, ?, ?)", + parsedSQL); + + Map<String, List<Integer>> legacyMap = new HashMap<>(); + String legacyParsed = FieldNamedPreparedStatement.parseNamedStatement(insertSQL, legacyMap); + // Keep documenting the old parse signature: the special-name fix is only enabled when the + // current field list is passed to the parser. + Assertions.assertTrue(legacyMap.containsKey("GLREG")); + Assertions.assertTrue(legacyMap.containsKey("MY")); + Assertions.assertFalse(legacyMap.containsKey("GLREG#")); + Assertions.assertTrue(legacyParsed.contains("?#")); + } + + @Test + void parseNamedStatementShouldNotSwallowNextClickHousePlaceholder() { + String sql = SqlUtils.getInsertIntoStatement("sink_table", new String[] {"A", "B"}); + String[] fieldNames = {"A, :B", "A", "B"}; + + Map<String, List<Integer>> parameterMap = new HashMap<>(); + String parsedSQL = + FieldNamedPreparedStatement.parseNamedStatement(sql, parameterMap, fieldNames); + + Assertions.assertEquals("INSERT INTO sink_table (\"A\", \"B\") VALUES (?, ?)", parsedSQL); + Assertions.assertFalse(parameterMap.containsKey("A, :B")); + Assertions.assertEquals("[1]", parameterMap.get("A").toString()); + Assertions.assertEquals("[2]", parameterMap.get("B").toString()); + } + + @Test + void prepareStatementBindsSpecialClickHouseColumnNames() throws Exception { + Connection connection = Mockito.mock(Connection.class); + PreparedStatement delegate = Mockito.mock(PreparedStatement.class); + String insertSQL = + SqlUtils.getInsertIntoStatement( + "sink_table", new String[] {"GLREG", "GLREG#", "MY COL", "COL-1"}); + + Mockito.when(connection.prepareStatement(Mockito.anyString())).thenReturn(delegate); + + PreparedStatement statement = + FieldNamedPreparedStatement.prepareStatement( + connection, insertSQL, new String[] {"GLREG", "GLREG#", "MY COL", "COL-1"}); + statement.setString(2, "region"); + statement.setString(3, "space"); + + Mockito.verify(connection) + .prepareStatement( + "INSERT INTO sink_table (\"GLREG\", \"GLREG#\", \"MY COL\", \"COL-1\") VALUES (?, ?, ?, ?)"); + Mockito.verify(delegate).setString(2, "region"); + Mockito.verify(delegate).setString(3, "space"); + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseIT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseIT.java index 47bd9672f9..87fe215152 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseIT.java +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/clickhouse/ClickhouseIT.java @@ -92,6 +92,8 @@ public class ClickhouseIT extends TestSuiteBase implements TestResource { private static final String SOURCE_TABLE = "source_table"; private static final String SOURCE_MERGE_TREE_TABLE = "source_merge_tree_table"; private static final String SINK_TABLE = "sink_table"; + private static final String SPECIAL_COLUMN_SOURCE_TABLE = "special_column_source_table"; + private static final String SPECIAL_COLUMN_SINK_TABLE = "special_column_sink_table"; private static final List<String> MULTI_SINK_TABLES = Arrays.asList("multi_sink_table1", "multi_sink_table2"); private static final List<String> MULTI_SOURCE_SINK_TABLES = @@ -123,6 +125,18 @@ public class ClickhouseIT extends TestSuiteBase implements TestResource { Assertions.assertEquals(0, execResult.getExitCode()); } + @TestTemplate + public void testClickhouseSinkWithSpecialCharactersInColumnNames(TestContainer container) + throws Exception { + initializeClickhouseSpecialColumnTable(); + Container.ExecResult execResult = + container.executeJob("/clickhouse_special_columns_to_clickhouse.conf"); + Assertions.assertEquals(0, execResult.getExitCode()); + assertSpecialColumnTableRows(); + dropTable(DATABASE + "." + SPECIAL_COLUMN_SOURCE_TABLE); + dropTable(DATABASE + "." + SPECIAL_COLUMN_SINK_TABLE); + } + @TestTemplate public void testClickhouseWithCreateSchemaWhenComment(TestContainer container) throws Exception { @@ -529,6 +543,71 @@ public class ClickhouseIT extends TestSuiteBase implements TestResource { } } + private void initializeClickhouseSpecialColumnTable() { + try { + Statement statement = this.connection.createStatement(); + String createSourceTable = + String.format( + "create table if not exists %s.%s(\n" + + " `id` Int64,\n" + + " `GLREG` String,\n" + + " `GLREG#` String,\n" + + " `MY COL` String,\n" + + " `COL-1` String\n" + + ")engine=MergeTree ORDER BY(id)", + DATABASE, SPECIAL_COLUMN_SOURCE_TABLE); + String createSinkTable = + String.format( + "create table if not exists %s.%s(\n" + + " `id` Int64,\n" + + " `GLREG` String,\n" + + " `GLREG#` String,\n" + + " `MY COL` String,\n" + + " `COL-1` String\n" + + ")engine=MergeTree ORDER BY(id)", + DATABASE, SPECIAL_COLUMN_SINK_TABLE); + statement.execute(createSourceTable); + statement.execute(createSinkTable); + statement.execute( + String.format( + "insert into %s.%s (`id`, `GLREG`, `GLREG#`, `MY COL`, `COL-1`)" + + " values (1, 'normal', 'hash', 'space', 'dash')," + + " (2, 'row2', 'row2#', 'row2 col', 'row2-1')," + + " (3, 'row3', 'row3#', 'row3 col', 'row3-1')", + DATABASE, SPECIAL_COLUMN_SOURCE_TABLE)); + } catch (SQLException e) { + throw new RuntimeException("Initializing Clickhouse table failed!", e); + } + } + + private void assertSpecialColumnTableRows() { + String expected = + "1,normal,hash,space,dash;" + + "2,row2,row2#,row2 col,row2-1;" + + "3,row3,row3#,row3 col,row3-1"; + try (Statement statement = this.connection.createStatement(); + ResultSet resultSet = + statement.executeQuery( + String.format( + "select `id`, `GLREG`, `GLREG#`, `MY COL`, `COL-1`" + + " from %s.%s order by `id`", + DATABASE, SPECIAL_COLUMN_SINK_TABLE))) { + StringBuilder actual = new StringBuilder(); + while (resultSet.next()) { + if (actual.length() > 0) { + actual.append(";"); + } + actual.append(resultSet.getLong(1)); + for (int i = 2; i <= 5; i++) { + actual.append(",").append(resultSet.getString(i)); + } + } + Assertions.assertEquals(expected, actual.toString()); + } catch (SQLException e) { + throw new RuntimeException("Querying Clickhouse table failed!", e); + } + } + private void initConnection() throws SQLException, ClassNotFoundException, InstantiationException, IllegalAccessException { diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/clickhouse_special_columns_to_clickhouse.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/clickhouse_special_columns_to_clickhouse.conf new file mode 100644 index 0000000000..3e3a1639ef --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-clickhouse-e2e/src/test/resources/clickhouse_special_columns_to_clickhouse.conf @@ -0,0 +1,50 @@ +# +# 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. +# +###### +###### This job copies a table whose column names contain characters that +###### are not valid Java identifier parts ('#', ' ', '-') to make sure the +###### ClickHouse sink named-parameter SQL still binds every field. The +###### column names come from the source table schema, so they never appear +###### as keys in this config file. +###### + +env { + parallelism = 1 + job.mode = "BATCH" + checkpoint.interval = 10 +} + +source { + Clickhouse { + host = "clickhouse:8123" + table_path = "default.special_column_source_table" + sql = "select * from special_column_source_table" + username = "default" + password = "" + plugin_output = "special_column_source_table" + } +} + +sink { + Clickhouse { + host = "clickhouse:8123" + database = "default" + table = "special_column_sink_table" + username = "default" + password = "" + } +}
