This is an automated email from the ASF dual-hosted git repository. JackieTien97 pushed a commit to branch rc/2.0.11 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit d5dc8f5492892d2638ee6c74ded5071da1eda146 Author: Caideyipi <[email protected]> AuthorDate: Mon Jul 20 15:24:25 2026 +0800 [Pipe] Reduce AirGap sink reconnect log noise (#18243) --- .../apache/iotdb/commons/i18n/PipeMessages.java | 8 ++-- .../apache/iotdb/commons/i18n/PipeMessages.java | 8 ++-- .../task/subtask/PipeAbstractSinkSubtask.java | 9 +++-- .../commons/pipe/resource/log/PipeLogger.java | 4 +- .../pipe/sink/protocol/IoTDBAirGapSink.java | 12 ++++-- .../commons/pipe/resource/PipeLoggerTest.java | 47 +++++++++------------- 6 files changed, 42 insertions(+), 46 deletions(-) diff --git a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java index cfeb392a3b5..02422dc4a8d 100644 --- a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java @@ -432,11 +432,11 @@ public final class PipeMessages { public static final String NON_CRITICAL_EXCEPTION_WILL_THROW_CRITICAL = "A non PipeRuntimeSinkCriticalException occurred, will throw a PipeRuntimeSinkCriticalException."; public static final String PIPE_CONNECTION_EXCEPTION_RETRYING = - "PipeConnectionException occurred, %s retries to handshake with the target system."; + "PipeConnectionException occurred, %s retries to handshake with the target system. Root cause: %s."; public static final String HANDSHAKE_SUCCESS = "{} handshakes with the target system successfully."; public static final String HANDSHAKE_FAILED_RETRYING = "{} failed to handshake with the target system for {} times, " - + "will retry at most {} times."; + + "will retry at most {} times. Root cause: {}."; public static final String INTERRUPTED_WHILE_SLEEPING_RETRY_HANDSHAKE = "Interrupted while sleeping, will retry to handshake with the target system."; public static final String HANDSHAKE_FAILED_STOPPING = @@ -511,9 +511,9 @@ public final class PipeMessages { public static final String FAILED_TO_CONNECT_TO_TARGET = "Failed to connect to target server ip: {}, port: {}, because: {}. Ignore it."; public static final String HANDSHAKE_ERROR_RECEIVING_END = - "Handshake error occurs. It may be caused by an error on the receiving end. Ignore it."; + "Handshake error occurs. It may be caused by an error on the receiving end. Ignore it. Root cause: {}."; public static final String HANDSHAKE_ERROR_WITH_TARGET = - "Handshake error with target server, socket: %s"; + "Handshake error with target server, endpoint: %s"; public static final String HANDSHAKE_SUCCESS_SOCKET = "Handshake success. Socket: {}"; public static final String FAILED_TO_CLOSE_CLIENT = "Failed to close client {}."; public static final String UNKNOWN_LOAD_BALANCE_STRATEGY = diff --git a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java index b1f5e59eeb1..ca5dc787125 100644 --- a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java @@ -415,10 +415,10 @@ public final class PipeMessages { public static final String NON_CRITICAL_EXCEPTION_WILL_THROW_CRITICAL = "发生非 PipeRuntimeSinkCriticalException,将抛出 PipeRuntimeSinkCriticalException。"; public static final String PIPE_CONNECTION_EXCEPTION_RETRYING = - "发生 PipeConnectionException,%s 正在重试与目标系统握手。"; + "发生 PipeConnectionException,%s 正在重试与目标系统握手。根因:%s。"; public static final String HANDSHAKE_SUCCESS = "{} 与目标系统握手成功。"; public static final String HANDSHAKE_FAILED_RETRYING = - "{} 与目标系统握手失败,已重试 {} 次,最多重试 {} 次。"; + "{} 与目标系统握手失败,已重试 {} 次,最多重试 {} 次。根因:{}。"; public static final String INTERRUPTED_WHILE_SLEEPING_RETRY_HANDSHAKE = "休眠时被中断,将重试与目标系统握手。"; public static final String HANDSHAKE_FAILED_STOPPING = @@ -491,9 +491,9 @@ public final class PipeMessages { public static final String FAILED_TO_CONNECT_TO_TARGET = "连接目标服务器失败(ip:{},端口:{}),原因:{}。已忽略。"; public static final String HANDSHAKE_ERROR_RECEIVING_END = - "握手错误,可能由接收端错误引起。已忽略。"; + "握手错误,可能由接收端错误引起。已忽略。根因:{}。"; public static final String HANDSHAKE_ERROR_WITH_TARGET = - "与目标服务器握手失败,socket:%s"; + "与目标服务器握手失败,endpoint:%s"; public static final String HANDSHAKE_SUCCESS_SOCKET = "握手成功。Socket:{}"; public static final String FAILED_TO_CLOSE_CLIENT = "关闭客户端 {} 失败。"; public static final String UNKNOWN_LOAD_BALANCE_STRATEGY = diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractSinkSubtask.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractSinkSubtask.java index 1b51e460956..ceaade1bb1f 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractSinkSubtask.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractSinkSubtask.java @@ -186,9 +186,9 @@ public abstract class PipeAbstractSinkSubtask extends PipeReportableSubtask { private boolean onPipeConnectionException(final Throwable throwable) { PipeLogger.log( LOGGER::warn, - throwable, PipeMessages.PIPE_CONNECTION_EXCEPTION_RETRYING, - outputPipeSink.getClass().getName()); + outputPipeSink.getClass().getName(), + ErrorHandlingCommonUtils.getRootCause(throwable).toString()); int retry = 0; while (retry < MAX_RETRY_TIMES) { @@ -200,12 +200,13 @@ public abstract class PipeAbstractSinkSubtask extends PipeReportableSubtask { break; } catch (final Exception e) { retry++; - LOGGER.warn( + PipeLogger.log( + LOGGER::warn, PipeMessages.HANDSHAKE_FAILED_RETRYING, outputPipeSink.getClass().getName(), retry, MAX_RETRY_TIMES, - e); + ErrorHandlingCommonUtils.getRootCause(e).toString()); try { sleepIfNoHighPriorityTask(getHandshakeRetrySleepInterval(e, retry)); } catch (final InterruptedException interruptedException) { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/log/PipeLogger.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/log/PipeLogger.java index 720999aae97..b906fc432cd 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/log/PipeLogger.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/log/PipeLogger.java @@ -26,9 +26,7 @@ import java.io.PrintStream; import java.util.function.Consumer; public class PipeLogger { - private static PipePeriodicalLogger logger = - (loggerFunction, rawMessage, formatter) -> - loggerFunction.accept(LoggerPeriodicalLogReducer.formatMessage(rawMessage, formatter)); + private static PipePeriodicalLogger logger = LoggerPeriodicalLogReducer::log; public static void log( final Consumer<String> loggerFunction, final String rawMessage, final Object... formatter) { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBAirGapSink.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBAirGapSink.java index 29a679e8037..8bcca0315f4 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBAirGapSink.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBAirGapSink.java @@ -23,8 +23,10 @@ import org.apache.iotdb.common.rpc.thrift.TEndPoint; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.i18n.PipeMessages; import org.apache.iotdb.commons.pipe.config.PipeConfig; +import org.apache.iotdb.commons.pipe.resource.log.PipeLogger; import org.apache.iotdb.commons.pipe.sink.payload.airgap.AirGapELanguageConstant; import org.apache.iotdb.commons.pipe.sink.payload.airgap.AirGapOneByteResponse; +import org.apache.iotdb.commons.utils.ErrorHandlingCommonUtils; import org.apache.iotdb.pipe.api.annotation.TableModel; import org.apache.iotdb.pipe.api.annotation.TreeModel; import org.apache.iotdb.pipe.api.customizer.configuration.PipeConnectorRuntimeConfiguration; @@ -181,7 +183,7 @@ public abstract class IoTDBAirGapSink extends IoTDBSink { socket.connect(new InetSocketAddress(ip, port), handshakeTimeoutMs); socket.setKeepAlive(true); sockets.set(i, socket); - LOGGER.info(PipeMessages.CONNECTED_TO_TARGET_SERVER, ip, port); + LOGGER.debug(PipeMessages.CONNECTED_TO_TARGET_SERVER, ip, port); failLogTimes.remove(nodeUrls.get(i)); } catch (final Exception e) { final TEndPoint endPoint = nodeUrls.get(i); @@ -198,7 +200,10 @@ public abstract class IoTDBAirGapSink extends IoTDBSink { sendHandshakeReq(socket); isSocketAlive.set(i, true); } catch (Exception e) { - LOGGER.warn(PipeMessages.HANDSHAKE_ERROR_RECEIVING_END, e); + PipeLogger.log( + LOGGER::warn, + PipeMessages.HANDSHAKE_ERROR_RECEIVING_END, + ErrorHandlingCommonUtils.getRootCause(e).toString()); } } @@ -219,7 +224,8 @@ public abstract class IoTDBAirGapSink extends IoTDBSink { if (!send(socket, generateHandShakeV2Payload())) { supportModsIfIsDataNodeReceiver = false; if (!send(socket, generateHandShakeV1Payload())) { - throw new PipeConnectionException(PipeMessages.HANDSHAKE_ERROR_WITH_TARGET + socket); + throw new PipeConnectionException( + String.format(PipeMessages.HANDSHAKE_ERROR_WITH_TARGET, socket.getEndPoint())); } } else { supportModsIfIsDataNodeReceiver = true; diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/resource/PipeLoggerTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/resource/PipeLoggerTest.java index d37c109c9f5..651bd552081 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/resource/PipeLoggerTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/resource/PipeLoggerTest.java @@ -24,20 +24,27 @@ import org.apache.iotdb.commons.pipe.resource.log.PipeLogger; import org.junit.Assert; import org.junit.Test; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; public class PipeLoggerTest { + @Test + public void testDefaultLoggerReducesDuplicateMessages() { + final AtomicInteger logCount = new AtomicInteger(0); + final String message = "PipeLoggerTest-" + System.nanoTime(); + + PipeLogger.log(ignored -> logCount.incrementAndGet(), message); + PipeLogger.log(ignored -> logCount.incrementAndGet(), message); + + Assert.assertEquals(1, logCount.get()); + } + @Test public void testLogMessageWithPercent() { final AtomicReference<String> message = new AtomicReference<>(); - setStringFormatLogger(); - try { - PipeLogger.log(message::set, "data_{sink.password=%/broken}"); - } finally { - setStringFormatLogger(); - } + PipeLogger.log(message::set, "data_{sink.password=%/broken}"); Assert.assertEquals("data_{sink.password=%/broken}", message.get()); } @@ -46,12 +53,7 @@ public class PipeLoggerTest { public void testLogMessageWithSlf4jPlaceholder() { final AtomicReference<String> message = new AtomicReference<>(); - setStringFormatLogger(); - try { - PipeLogger.log(message::set, "PipeLoggerCacheMaxSizeInBytes: {}", 1024); - } finally { - setStringFormatLogger(); - } + PipeLogger.log(message::set, "PipeLoggerCacheMaxSizeInBytes: {}", 1024); Assert.assertEquals("PipeLoggerCacheMaxSizeInBytes: 1024", message.get()); } @@ -60,24 +62,13 @@ public class PipeLoggerTest { public void testLogThrowableWithPercentInStackTrace() { final AtomicReference<String> message = new AtomicReference<>(); - setStringFormatLogger(); - try { - PipeLogger.log( - message::set, - new RuntimeException("data_{sink.password=%/broken}"), - "Failed to transfer event %s", - "root.sg.d1"); - } finally { - setStringFormatLogger(); - } + PipeLogger.log( + message::set, + new RuntimeException("data_{sink.password=%/broken}"), + "Failed to transfer event %s", + "root.sg.d1"); Assert.assertTrue(message.get().contains("Failed to transfer event root.sg.d1")); Assert.assertTrue(message.get().contains("data_{sink.password=%/broken}")); } - - private void setStringFormatLogger() { - PipeLogger.setLogger( - (loggerFunction, rawMessage, formatter) -> - loggerFunction.accept(String.format(rawMessage, formatter))); - } }
