This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 339fb88cce [Improve][Core] Improve config log desensitization coverage
(#11811)
339fb88cce is described below
commit 339fb88cce0d9a5e158650b331de405673002aa2
Author: Jast <[email protected]>
AuthorDate: Tue Aug 25 22:45:13 2026 +0800
[Improve][Core] Improve config log desensitization coverage (#11811)
Co-authored-by: David Zollo <[email protected]>
---
.../configuration/config-encryption-decryption.md | 6 +-
.../configuration/config-encryption-decryption.md | 5 +-
.../core/starter/utils/ConfigBuilder.java | 78 ++++++++++++++--
.../core/starter/utils/ConfigShadeUtils.java | 16 ++++
.../core/starter/utils/ConfigBuilderTest.java | 103 ++++++++++++++++++++-
.../core/starter/utils/ConfigShadeTest.java | 33 +++++++
6 files changed, 230 insertions(+), 11 deletions(-)
diff --git a/docs/en/introduction/configuration/config-encryption-decryption.md
b/docs/en/introduction/configuration/config-encryption-decryption.md
index b6b021fe8b..64327bae38 100644
--- a/docs/en/introduction/configuration/config-encryption-decryption.md
+++ b/docs/en/introduction/configuration/config-encryption-decryption.md
@@ -18,6 +18,10 @@ Base64 encryption support encrypt the following parameters
by default:
And users can add custom parameters to `shade.options` for encryption and
decryption.
+When SeaTunnel prints parsed configuration to logs, it also masks these
options in nested
+configuration paths. The log masking matcher treats `.`, `_`, and `-` as
equivalent separators, so
+for example `access.key`, `access_key`, and `access-key` can share the same
masking rule.
+
Next, I'll show how to quickly use SeaTunnel's own `base64` encryption:
1. And new option `shade.identifier` and `shade.options` in env block of
config file, `shade.identifier` indicate what the encryption method that you
want to use, while `shade.options` specifies which parameters should be
encrypted/decrypted. In this example, we should add `shade.identifier = base64`
in config as the following shown:
@@ -231,4 +235,4 @@ If you want to encrypt and decrypt with customized params,
you can follow the st
public String decrypt(String content) {
return content.substring(0, content.length() - suffix.length());
}
- ```
\ No newline at end of file
+ ```
diff --git a/docs/zh/introduction/configuration/config-encryption-decryption.md
b/docs/zh/introduction/configuration/config-encryption-decryption.md
index fafc4bf8ae..de795d3741 100644
--- a/docs/zh/introduction/configuration/config-encryption-decryption.md
+++ b/docs/zh/introduction/configuration/config-encryption-decryption.md
@@ -18,6 +18,9 @@ Base64编码默认支持加密以下参数:
用户也可以在 `shade.options` 指定要用于加解密的参数.
+当 SeaTunnel 将解析后的配置打印到日志时,也会在嵌套配置路径中对这些参数做脱敏展示。日志脱敏匹配会将
+`.`、`_` 和 `-` 视为等价分隔符,因此 `access.key`、`access_key` 和 `access-key` 可以复用同一条脱敏规则。
+
接下来,将展示如何快速使用 SeaTunnel 自带的 `base64` 加密功能:
1. 在配置文件的环境变量(env)部分新增了选项 `shade.identifier` 和
`shade.options`。`shade.identifier`用于表示您想要使用的加密方法,`shade.options`用于指定您想加解密的参数。
@@ -232,4 +235,4 @@ Base64编码默认支持加密以下参数:
public String decrypt(String content) {
return content.substring(0, content.length() - suffix.length());
}
- ```
\ No newline at end of file
+ ```
diff --git
a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigBuilder.java
b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigBuilder.java
index 6f0986f94e..7d409c60e0 100644
---
a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigBuilder.java
+++
b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigBuilder.java
@@ -57,6 +57,10 @@ public class ConfigBuilder {
ConfigRenderOptions.concise().setFormatted(true);
private static final String PLACEHOLDER_REGEX =
"\\$\\{([^:{}]+)(?::[^}]*)?\\}";
+ private static final String MASKED_VALUE = "******";
+ private static final String CONFIG_PATH_SEPARATOR = ".";
+ // Treat common option separators as equivalent when matching config paths
in logs.
+ private static final Pattern CONFIG_OPTION_SEPARATOR_PATTERN =
Pattern.compile("[._-]+");
private ConfigBuilder() {
// utility class and cannot be instantiated
@@ -95,7 +99,7 @@ public class ConfigBuilder {
mapToString(
configDesensitization(
config.root().unwrapped(),
-
ConfigShadeUtils.getSensitiveOptions(config))));
+
ConfigShadeUtils.getLogDesensitizationOptions(config))));
return config;
}
@@ -112,40 +116,62 @@ public class ConfigBuilder {
mapToString(
configDesensitization(
config.root().unwrapped(),
-
ConfigShadeUtils.getSensitiveOptions(config))));
+
ConfigShadeUtils.getLogDesensitizationOptions(config))));
return config;
}
public static Map<String, Object> configDesensitization(
Map<String, Object> configMap, Set<String> sensitiveKeywords) {
+ Set<String> normalizedSensitiveKeywords =
+ sensitiveKeywords.stream()
+ .map(ConfigBuilder::normalizeConfigOption)
+ .collect(Collectors.toSet());
+ return configDesensitization(configMap, normalizedSensitiveKeywords,
null);
+ }
+
+ /**
+ * Recursively builds a masked copy of the config map.
+ *
+ * <p>The accumulated {@code parentPath} preserves dotted option context
after HOCON has
+ * expanded paths into nested maps.
+ */
+ private static Map<String, Object> configDesensitization(
+ Map<String, Object> configMap,
+ Set<String> normalizedSensitiveKeywords,
+ String parentPath) {
return configMap.entrySet().stream()
.collect(
LinkedHashMap::new,
(m, p) -> {
String key = p.getKey();
Object value = p.getValue();
- if (sensitiveKeywords.contains(key.toLowerCase()))
{
+ String configPath =
+ parentPath == null
+ ? key
+ : parentPath +
CONFIG_PATH_SEPARATOR + key;
+ if (isSensitiveOption(key, configPath,
normalizedSensitiveKeywords)) {
if (value instanceof List<?>) {
List<Object> maskedList =
((List<?>) value)
.stream()
- .map(v -> "******")
+ .map(v ->
MASKED_VALUE)
.collect(Collectors.toList());
m.put(key, maskedList);
} else {
- m.put(key, "******");
+ m.put(key, MASKED_VALUE);
}
} else if (value instanceof String
&& ((String) value)
.regionMatches(true, 0, "jdbc:",
0, "jdbc:".length())) {
- m.put(key, "******");
+ m.put(key, MASKED_VALUE);
} else {
if (value instanceof Map<?, ?>) {
m.put(
key,
configDesensitization(
(Map<String, Object>)
value,
- sensitiveKeywords));
+
normalizedSensitiveKeywords,
+ configPath));
} else if (value instanceof List<?>) {
List<?> listValue = (List<?>) value;
List<Object> newList =
@@ -155,7 +181,8 @@ public class ConfigBuilder {
if (v
instanceof Map<?, ?>) {
return
configDesensitization(
(Map<String, Object>) v,
-
sensitiveKeywords);
+
normalizedSensitiveKeywords,
+
configPath);
} else {
return v;
}
@@ -170,6 +197,41 @@ public class ConfigBuilder {
LinkedHashMap::putAll);
}
+ /**
+ * Checks whether the current option should be masked in the parsed-config
log.
+ *
+ * <p>The matcher compares both the leaf key and the accumulated config
path. Option separators
+ * '.', '_' and '-' are treated as equivalent, so paths like {@code
+ * kafka.config.sasl.jaas.config} can match {@code sasl.jaas.config}.
Suffix matching is applied
+ * only to multi-segment sensitive options such as {@code access_key};
single-word options such
+ * as {@code token} still require an exact leaf-key or full-path match.
+ */
+ private static boolean isSensitiveOption(
+ String key, String configPath, Set<String>
normalizedSensitiveKeywords) {
+ String normalizedKey = normalizeConfigOption(key);
+ String normalizedConfigPath = normalizeConfigOption(configPath);
+ if (normalizedSensitiveKeywords.contains(normalizedKey)
+ || normalizedSensitiveKeywords.contains(normalizedConfigPath))
{
+ return true;
+ }
+ return normalizedSensitiveKeywords.stream()
+ .filter(ConfigBuilder::isMultiSegmentOption)
+ .anyMatch(
+ sensitiveKeyword -> normalizedConfigPath.endsWith("_"
+ sensitiveKeyword));
+ }
+
+ /**
+ * Normalizes common option separator styles so equivalent config names
can share one matching
+ * rule.
+ */
+ private static String normalizeConfigOption(String option) {
+ return
CONFIG_OPTION_SEPARATOR_PATTERN.matcher(option.toLowerCase()).replaceAll("_");
+ }
+
+ private static boolean isMultiSegmentOption(String option) {
+ return option.contains("_");
+ }
+
public static Config of(
@NonNull ConfigAdapter configAdapter, @NonNull Path filePath,
List<String> variables) {
log.info("With config adapter spi {}",
configAdapter.getClass().getName());
diff --git
a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java
b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java
index 92b3096e4c..40ce48fcce 100644
---
a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java
+++
b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigShadeUtils.java
@@ -54,6 +54,9 @@ public final class ConfigShadeUtils {
public static final String[] DEFAULT_SENSITIVE_KEYWORDS =
new String[] {"password", "username", "auth", "token",
"access_key", "secret_key"};
+ private static final String[] DEFAULT_LOG_MASK_ONLY_KEYWORDS =
+ new String[] {"sasl.jaas.config"};
+
private static final Map<String, ConfigShade> CONFIG_SHADES = new
HashMap<>();
private static final ConfigShade DEFAULT_SHADE = new DefaultConfigShade();
@@ -240,6 +243,19 @@ public final class ConfigShadeUtils {
return sensitiveOptions;
}
+ /**
+ * Returns option names used only for parsed-config log masking.
+ *
+ * <p>This method extends the encryption/decryption option list with
log-only option names.
+ * Adding entries here changes only the rendered log output and does not
change how existing
+ * configs are encrypted or decrypted.
+ */
+ public static Set<String> getLogDesensitizationOptions(Config config) {
+ Set<String> sensitiveOptions = getSensitiveOptions(config);
+ sensitiveOptions.addAll(Arrays.asList(DEFAULT_LOG_MASK_ONLY_KEYWORDS));
+ return sensitiveOptions;
+ }
+
public static class Base64ConfigShade implements ConfigShade {
private static final Base64.Encoder ENCODER = Base64.getEncoder();
diff --git
a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java
b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java
index ee96376de7..43f711de3c 100644
---
a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java
+++
b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigBuilderTest.java
@@ -58,7 +58,7 @@ public class ConfigBuilderTest {
Map<String, Object> desensitized =
ConfigBuilder.configDesensitization(
- config, ConfigShadeUtils.getSensitiveOptions(null));
+ config,
ConfigShadeUtils.getLogDesensitizationOptions(null));
List<?> sources = (List<?>) desensitized.get("source");
Map<?, ?> desensitizedSource = (Map<?, ?>) sources.get(0);
@@ -66,4 +66,105 @@ public class ConfigBuilderTest {
Assertions.assertEquals(
"https://catalog.example.com/tables",
desensitizedSource.get("metadata_url"));
}
+
+ @Test
+ public void testConfigDesensitizationMasksKafkaJaasConfig() {
+ Map<String, Object> jaasConfig = new LinkedHashMap<>();
+ jaasConfig.put(
+ "config",
+ "org.apache.kafka.common.security.scram.ScramLoginModule
required "
+ + "username=\"alice\" password=\"secret\";");
+
+ Map<String, Object> saslConfig = new LinkedHashMap<>();
+ saslConfig.put("jaas", jaasConfig);
+
+ Map<String, Object> kafkaConfig = new LinkedHashMap<>();
+ kafkaConfig.put("bootstrap.servers", "localhost:9092");
+ kafkaConfig.put("sasl", saslConfig);
+
+ Map<String, Object> source = new LinkedHashMap<>();
+ source.put("kafka.config", kafkaConfig);
+
+ Map<String, Object> config = new LinkedHashMap<>();
+ config.put("source", Arrays.asList(source));
+
+ Map<String, Object> desensitized =
+ ConfigBuilder.configDesensitization(
+ config,
ConfigShadeUtils.getLogDesensitizationOptions(null));
+ List<?> sources = (List<?>) desensitized.get("source");
+ Map<?, ?> desensitizedSource = (Map<?, ?>) sources.get(0);
+ Map<?, ?> desensitizedKafkaConfig = (Map<?, ?>)
desensitizedSource.get("kafka.config");
+ Map<?, ?> desensitizedSaslConfig = (Map<?, ?>)
desensitizedKafkaConfig.get("sasl");
+ Map<?, ?> desensitizedJaasConfig = (Map<?, ?>)
desensitizedSaslConfig.get("jaas");
+
+ Assertions.assertEquals("******",
desensitizedJaasConfig.get("config"));
+ Assertions.assertEquals("localhost:9092",
desensitizedKafkaConfig.get("bootstrap.servers"));
+ }
+
+ @Test
+ public void testConfigDesensitizationMasksS3CredentialOptions() {
+ Map<String, Object> accessKeyConfig = new LinkedHashMap<>();
+ accessKeyConfig.put("key", "access-key");
+
+ Map<String, Object> s3aConfig = new LinkedHashMap<>();
+ s3aConfig.put("endpoint", "http://localhost:9000");
+ s3aConfig.put("path.style.access", "true");
+ s3aConfig.put("aws.credentials.provider",
"SimpleAWSCredentialsProvider");
+ s3aConfig.put("access", accessKeyConfig);
+
+ Map<String, Object> fsConfig = new LinkedHashMap<>();
+ fsConfig.put("s3a", s3aConfig);
+
+ Map<String, Object> checkpointConfig = new LinkedHashMap<>();
+ checkpointConfig.put("fs", fsConfig);
+ checkpointConfig.put("fs.s3a.access-key", "access-key");
+ checkpointConfig.put("fs.s3a.secret.key", "secret-key");
+
+ Map<String, Object> config = new LinkedHashMap<>();
+ config.put("checkpoint", checkpointConfig);
+
+ Map<String, Object> desensitized =
+ ConfigBuilder.configDesensitization(
+ config, ConfigShadeUtils.getSensitiveOptions(null));
+ Map<?, ?> desensitizedCheckpoint = (Map<?, ?>)
desensitized.get("checkpoint");
+ Map<?, ?> desensitizedFsConfig = (Map<?, ?>)
desensitizedCheckpoint.get("fs");
+ Map<?, ?> desensitizedS3aConfig = (Map<?, ?>)
desensitizedFsConfig.get("s3a");
+ Map<?, ?> desensitizedAccessConfig = (Map<?, ?>)
desensitizedS3aConfig.get("access");
+
+ Assertions.assertEquals("******", desensitizedAccessConfig.get("key"));
+ Assertions.assertEquals("******",
desensitizedCheckpoint.get("fs.s3a.access-key"));
+ Assertions.assertEquals("******",
desensitizedCheckpoint.get("fs.s3a.secret.key"));
+ Assertions.assertEquals("http://localhost:9000",
desensitizedS3aConfig.get("endpoint"));
+ Assertions.assertEquals("true",
desensitizedS3aConfig.get("path.style.access"));
+ Assertions.assertEquals(
+ "SimpleAWSCredentialsProvider",
+ desensitizedS3aConfig.get("aws.credentials.provider"));
+ }
+
+ @Test
+ public void testConfigDesensitizationKeepsSchemaFieldsReadable() {
+ Map<String, Object> fields = new LinkedHashMap<>();
+ fields.put("access_token", "string");
+ fields.put("user-password", "string");
+
+ Map<String, Object> schema = new LinkedHashMap<>();
+ schema.put("fields", fields);
+
+ Map<String, Object> source = new LinkedHashMap<>();
+ source.put("schema", schema);
+
+ Map<String, Object> config = new LinkedHashMap<>();
+ config.put("source", Arrays.asList(source));
+
+ Map<String, Object> desensitized =
+ ConfigBuilder.configDesensitization(
+ config,
ConfigShadeUtils.getLogDesensitizationOptions(null));
+ List<?> sources = (List<?>) desensitized.get("source");
+ Map<?, ?> desensitizedSource = (Map<?, ?>) sources.get(0);
+ Map<?, ?> desensitizedSchema = (Map<?, ?>)
desensitizedSource.get("schema");
+ Map<?, ?> desensitizedFields = (Map<?, ?>)
desensitizedSchema.get("fields");
+
+ Assertions.assertEquals("string",
desensitizedFields.get("access_token"));
+ Assertions.assertEquals("string",
desensitizedFields.get("user-password"));
+ }
}
diff --git
a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigShadeTest.java
b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigShadeTest.java
index 8e1a1f48d1..c65471a7dd 100644
---
a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigShadeTest.java
+++
b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/utils/ConfigShadeTest.java
@@ -41,6 +41,7 @@ import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Base64;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -82,6 +83,38 @@ public class ConfigShadeTest {
config.getConfigList("source").get(0).getString("secret_key"),
SECRET_KEY);
}
+ @Test
+ public void testLogOnlySensitiveOptionsDoNotChangeDefaultDecryptOptions() {
+ String jaasConfig =
+ "org.apache.kafka.common.security.scram.ScramLoginModule
required "
+ + "username=\"alice\" password=\"plain\";";
+
+ Map<String, Object> env = new LinkedHashMap<>();
+ env.put("shade.identifier", "base64");
+
+ Map<String, Object> source = new LinkedHashMap<>();
+ source.put("plugin_name", "FakeSource");
+ source.put(
+ "username",
+
Base64.getEncoder().encodeToString(USERNAME.getBytes(StandardCharsets.UTF_8)));
+ source.put("sasl.jaas.config", jaasConfig);
+
+ Map<String, Object> sink = new LinkedHashMap<>();
+ sink.put("plugin_name", "Console");
+
+ Map<String, Object> configMap = new LinkedHashMap<>();
+ configMap.put("env", env);
+ configMap.put("source", Arrays.asList(source));
+ configMap.put("sink", Arrays.asList(sink));
+
+ Config config =
ConfigShadeUtils.decryptConfig(ConfigFactory.parseMap(configMap));
+
+ Assertions.assertEquals(
+ USERNAME,
config.getConfigList("source").get(0).getString("username"));
+ Assertions.assertEquals(
+ jaasConfig,
config.getConfigList("source").get(0).getString("sasl.jaas.config"));
+ }
+
@Test
public void testUsePrivacyHandlerHocon() throws URISyntaxException {
URL resource = ConfigShadeTest.class.getResource("/config.shade.conf");