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())));
+ }
+ }
+}