This is an automated email from the ASF dual-hosted git repository.

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new 0a60d39e2d fix(logging): support standalone ClickHouse tables (#7262)
0a60d39e2d is described below

commit 0a60d39e2dbdaca00dfb32f6853749f93f0b98fb
Author: Liming Deng <[email protected]>
AuthorDate: Sun Sep 27 17:43:43 2026 +0800

    fix(logging): support standalone ClickHouse tables (#7262)
    
    * fix(logging): support standalone ClickHouse tables
    
    * fix(logging): make ClickHouse table routing and lifecycle explicit
    
    ---------
    
    Co-authored-by: aias00 <[email protected]>
---
 .../client/ClickHouseLogCollectClient.java         |  15 ++-
 .../constant/ClickHouseLoggingConstant.java        |  15 ++-
 .../client/ClickHouseTableRoutingTest.java         | 107 +++++++++++++++++++++
 3 files changed, 131 insertions(+), 6 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseLogCollectClient.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseLogCollectClient.java
index 823ae373c1..8ef5cb56bf 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseLogCollectClient.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseLogCollectClient.java
@@ -52,6 +52,8 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
 
     private String database;
 
+    private String insertSql;
+
     /**
      * consume logs.
      * @param logs logs
@@ -60,6 +62,9 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
     @Override
     public void consume0(@NonNull final List<ShenyuRequestLog> logs) throws 
Exception {
         if (CollectionUtils.isNotEmpty(logs)) {
+            if (Objects.isNull(insertSql)) {
+                throw new IllegalStateException("ClickHouse log client must be 
initialized successfully before consuming logs");
+            }
             Object[][] datas = new Object[logs.size()][];
             for (int i = 0; i < logs.size(); i++) {
                 Object[] data = new Object[] {
@@ -85,7 +90,7 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
                 };
                 datas[i] = data;
             }
-            ClickHouseClient.send(endpoint, 
String.format(ClickHouseLoggingConstant.PRE_INSERT_SQL, database),
+            ClickHouseClient.send(endpoint, insertSql,
                     new ClickHouseValue[]{
                             ClickHouseOffsetDateTimeValue.ofNull(3, 
TimeZone.getTimeZone("Asia/Shanghai")),
                             ClickHouseStringValue.ofNull(),
@@ -112,6 +117,7 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
 
     @Override
     public void close0() {
+        insertSql = null;
         if (Objects.nonNull(client)) {
             client.close();
         }
@@ -129,6 +135,8 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
         final String password = config.getPassword();
         final String ttl = StringUtils.defaultIfBlank(config.getTtl(), "30");
         database = config.getDatabase();
+        boolean distributed = StringUtils.isNotBlank(config.getClusterName());
+        insertSql = null;
         endpoint = ClickHouseNode.builder()
             .host(config.getHost())
             .port(ClickHouseProtocol.HTTP, Integer.valueOf(config.getPort()))
@@ -139,12 +147,15 @@ public class ClickHouseLogCollectClient extends 
AbstractLogConsumeClient<ClickHo
             ClickHouseRequest<?> request = 
client.connect(endpoint).format(ClickHouseFormat.TabSeparatedWithNamesAndTypes);
             
request.query(String.format(ClickHouseLoggingConstant.CREATE_DATABASE_SQL, 
database)).executeAndWait();
             
request.query(String.format(ClickHouseLoggingConstant.CREATE_TABLE_SQL, 
database, config.getEngine(), ttl)).executeAndWait();
-            
request.query(String.format(ClickHouseLoggingConstant.CREATE_DISTRIBUTED_TABLE_SQL,
 database, database, config.getClusterName(), database)).executeAndWait();
+            if (distributed) {
+                
request.query(String.format(ClickHouseLoggingConstant.CREATE_DISTRIBUTED_TABLE_SQL,
 database, database, config.getClusterName(), database)).executeAndWait();
+            }
         } catch (Exception e) {
             LOG.error("inti ClickHouseLogClient error", e);
             close0();
             return false;
         }
+        insertSql = String.format(distributed ? 
ClickHouseLoggingConstant.PRE_INSERT_SQL : 
ClickHouseLoggingConstant.LOCAL_PRE_INSERT_SQL, database);
         return true;
     }
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/constant/ClickHouseLoggingConstant.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/constant/ClickHouseLoggingConstant.java
index 6ac9cdfd0f..7a282b07dd 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/constant/ClickHouseLoggingConstant.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/main/java/org/apache/shenyu/plugin/logging/clickhouse/constant/ClickHouseLoggingConstant.java
@@ -62,14 +62,21 @@ public class ClickHouseLoggingConstant {
     public static final String CREATE_DISTRIBUTED_TABLE_SQL = "create table if 
not exists `%s`.request_log_distributed\n"
             + " AS `%s`.request_log ENGINE = Distributed('%s', '%s', 
'request_log', rand());";
 
-    /**
-     * The constant PRE_INSERT_SQL.
-     */
-    public static final String PRE_INSERT_SQL = "INSERT INTO 
`%s`.request_log_distributed"
+    private static final String INSERT_SQL_TEMPLATE = "INSERT INTO `%s`.%s"
             + "(timeLocal, clientIp, method, requestHeader, responseHeader, 
queryParams, "
             + "requestBody, requestUri, responseBody, responseContentLength, 
rpcType, status, upstreamIp, upstreamResponseTime, userAgent, host, module, 
traceId, path) "
             + "VALUES "
             + "(:timeLocal, :clientIp,:method, :requestHeader, 
:responseHeader, :queryParams,"
             + " :requestBody, :requestUri, :responseBody, 
:responseContentLength, :rpcType, :status, :upstreamIp, :upstreamResponseTime, 
:userAgent, :host, :module, :traceId, :path);";
 
+    /**
+     * Insert logs through the cluster's distributed table.
+     */
+    public static final String PRE_INSERT_SQL = 
String.format(INSERT_SQL_TEMPLATE, "%s", "request_log_distributed");
+
+    /**
+     * Insert logs directly into a standalone server's local table.
+     */
+    public static final String LOCAL_PRE_INSERT_SQL = 
String.format(INSERT_SQL_TEMPLATE, "%s", "request_log");
+
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/test/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseTableRoutingTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/test/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseTableRoutingTest.java
new file mode 100644
index 0000000000..d8cac5b397
--- /dev/null
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-clickhouse/src/test/java/org/apache/shenyu/plugin/logging/clickhouse/client/ClickHouseTableRoutingTest.java
@@ -0,0 +1,107 @@
+/*
+ * 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.shenyu.plugin.logging.clickhouse.client;
+
+import com.clickhouse.client.ClickHouseClient;
+import com.clickhouse.client.ClickHouseClientBuilder;
+import com.clickhouse.client.ClickHouseNode;
+import com.clickhouse.client.ClickHouseRequest;
+import com.clickhouse.client.ClickHouseValue;
+import 
org.apache.shenyu.plugin.logging.clickhouse.config.ClickHouseLogCollectConfig.ClickHouseLogConfig;
+import 
org.apache.shenyu.plugin.logging.clickhouse.constant.ClickHouseLoggingConstant;
+import org.apache.shenyu.plugin.logging.common.entity.ShenyuRequestLog;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.NullAndEmptySource;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.MockedStatic;
+
+import java.util.Collections;
+import java.util.concurrent.CompletableFuture;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.RETURNS_SELF;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class ClickHouseTableRoutingTest {
+
+    @Test
+    void consumingBeforeInitializationReportsLifecycleError() {
+        IllegalStateException error = assertThrows(IllegalStateException.class,
+                () -> new 
ClickHouseLogCollectClient().consume0(Collections.singletonList(new 
ShenyuRequestLog())));
+        assertEquals("ClickHouse log client must be initialized successfully 
before consuming logs", error.getMessage());
+    }
+
+    @ParameterizedTest
+    @NullAndEmptySource
+    @ValueSource(strings = {" ", "logs-cluster"})
+    void routesDdlAndInsertsAccordingToClusterConfiguration(final String 
cluster) throws Exception {
+        ClickHouseLogConfig config = new ClickHouseLogConfig();
+        config.setHost("localhost");
+        config.setPort("8123");
+        config.setDatabase("logs");
+        config.setUsername("default");
+        config.setPassword("");
+        config.setEngine("MergeTree");
+        config.setClusterName(cluster);
+        ClickHouseClient client = mock(ClickHouseClient.class);
+        ClickHouseClientBuilder builder = mock(ClickHouseClientBuilder.class);
+        ClickHouseRequest<?> request = mock(ClickHouseRequest.class, 
RETURNS_SELF);
+        when(builder.build()).thenReturn(client);
+        doReturn(request).when(client).connect(any(ClickHouseNode.class));
+        boolean distributed = "logs-cluster".equals(cluster);
+        String ddl = 
String.format(ClickHouseLoggingConstant.CREATE_DISTRIBUTED_TABLE_SQL, "logs", 
"logs", cluster, "logs");
+        String insert = String.format(distributed ? 
ClickHouseLoggingConstant.PRE_INSERT_SQL : 
ClickHouseLoggingConstant.LOCAL_PRE_INSERT_SQL, "logs");
+        assertTrue(insert.startsWith("INSERT INTO `logs`." + (distributed ? 
"request_log_distributed" : "request_log") + "("));
+        try (MockedStatic<ClickHouseClient> clients = 
mockStatic(ClickHouseClient.class)) {
+            clients.when(ClickHouseClient::builder).thenReturn(builder);
+            clients.when(() -> 
ClickHouseClient.send(any(ClickHouseNode.class), anyString(), 
any(ClickHouseValue[].class), any(Object[][].class)))
+                    
.thenReturn(CompletableFuture.completedFuture(Collections.emptyList()));
+            ClickHouseLogCollectClient collector = new 
ClickHouseLogCollectClient();
+            try {
+                assertTrue(collector.initClient0(config));
+                
verify(request).query(String.format(ClickHouseLoggingConstant.CREATE_DATABASE_SQL,
 "logs"));
+                
verify(request).query(String.format(ClickHouseLoggingConstant.CREATE_TABLE_SQL, 
"logs", "MergeTree", "30"));
+                if (distributed) {
+                    verify(request).query(ddl);
+                } else {
+                    verify(request, never()).query(ddl);
+                }
+                ShenyuRequestLog log = new ShenyuRequestLog();
+                log.setTimeLocal("2026-09-25 00:00:00.123");
+                collector.consume0(Collections.singletonList(log));
+                clients.verify(() -> 
ClickHouseClient.send(any(ClickHouseNode.class), eq(insert), 
any(ClickHouseValue[].class), any(Object[][].class)));
+            } finally {
+                collector.close0();
+            }
+            verify(client).close();
+            assertThrows(IllegalStateException.class, () -> 
collector.consume0(Collections.singletonList(new ShenyuRequestLog())));
+        }
+    }
+}

Reply via email to