Skip to content

Commit a51d807

Browse files
authored
feat: Resolve env var placeholders in Flink config.yaml (#369)
1 parent ebf2382 commit a51d807

2 files changed

Lines changed: 39 additions & 1 deletion

File tree

flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/BaseRunner.java

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,7 @@ Configuration initConfiguration() {
9191
var conf = new Configuration();
9292
if (StringUtils.isNotBlank(configDir)) {
9393
log.info("Loading Flink configuration from '{}'", configDir);
94-
conf = GlobalConfiguration.loadConfiguration(configDir);
94+
conf = resolveEnvVars(GlobalConfiguration.loadConfiguration(configDir));
9595
}
9696

9797
// Do not overwrite runtime given in YAML
@@ -102,6 +102,20 @@ Configuration initConfiguration() {
102102
return conf;
103103
}
104104

105+
private Configuration resolveEnvVars(Configuration source) {
106+
var conf = new Configuration(source);
107+
source
108+
.toMap()
109+
.forEach(
110+
(key, value) -> {
111+
var resolved = resolver.resolve(value);
112+
if (!resolved.equals(value)) {
113+
conf.setString(key, resolved);
114+
}
115+
});
116+
return conf;
117+
}
118+
105119
static String readTextFile(String path) throws IOException {
106120
var f = new File(path);
107121
if (!f.exists()) {

flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerTest.java

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import java.io.IOException;
2727
import java.nio.file.Files;
2828
import java.nio.file.Path;
29+
import java.util.Map;
2930
import java.util.function.Supplier;
3031
import org.apache.flink.api.common.RuntimeExecutionMode;
3132
import org.apache.flink.configuration.ExecutionOptions;
@@ -171,4 +172,27 @@ void initConfiguration_shouldCreateNewConfigIfNoPathGiven(RuntimeExecutionMode m
171172
assertThat(finalConf.keySet()).hasSize(1);
172173
assertThat(finalConf.get(ExecutionOptions.RUNTIME_MODE)).isEqualTo(mode);
173174
}
175+
176+
@Test
177+
void initConfiguration_shouldResolveEnvVarPlaceholdersInConfig() throws IOException {
178+
// Arrange — a secret injected via env var, e.g. the Iceberg maintenance lock password
179+
Path configFile = tempDir.resolve("config.yaml");
180+
Files.writeString(
181+
configFile,
182+
"flink-maintenance.lock.jdbc.password: ${ICEBERG_LOCK_PASSWORD}\ndummy.key: literal");
183+
184+
var resolver =
185+
EnvVarResolver.builder().envVars(Map.of("ICEBERG_LOCK_PASSWORD", "s3cr3t")).build();
186+
var cliRunner =
187+
new CliRunner(
188+
RuntimeExecutionMode.STREAMING, resolver, null, null, tempDir.toString(), null);
189+
190+
// Act
191+
var finalConf = cliRunner.initConfiguration();
192+
193+
// Assert
194+
assertThat(finalConf.toMap())
195+
.containsEntry("flink-maintenance.lock.jdbc.password", "s3cr3t")
196+
.containsEntry("dummy.key", "literal");
197+
}
174198
}

0 commit comments

Comments
 (0)