Flink通过REST API提交作业时如何传递并获取自定义参数
方案可行性判断
你提出的「本地完成Redis密码解密,再通过REST API提交Flink作业时传入解密后密码」的方案完全可行,没有框架层面的限制。
无法获取参数的原因
你之前从StreamExecutionEnvironment.getExecutionEnvironment().getConfig()里查找参数的方向是错的:REST API提交时传入的programArgs、programArgsList参数,不会自动写入ExecutionEnvironment的配置中,这两类参数的本质是作业main方法的启动入参,和你本地IDE运行作业时传入的program arguments完全等价。
注意:
programArgs会把传入的字符串按空格拆分后生成main方法的args数组,适合无空格的简单参数;programArgsList会直接把列表元素按顺序映射为args数组,不会做额外拆分,推荐优先用这个字段传参,避免参数值含空格时解析出错。
参数获取的具体实现
直接在作业入口类的main方法中解析入参String[] args即可,推荐用Flink自带的ParameterTool做参数解析,步骤如下:
- 入口类main方法直接从args解析参数
public static void main(String[] args) throws Exception { // 解析REST传入的所有启动参数 ParameterTool jobParams = ParameterTool.fromArgs(args); // 读取解密后的Redis密码,缺失会直接抛异常 String redisPassword = jobParams.getRequired("redis.password"); // 初始化执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 如果需要在各个算子中使用参数,可注册为全局作业参数 env.getConfig().setGlobalJobParameters(jobParams); // 后续初始化Redis连接、编写作业业务逻辑即可 env.execute("RedisConnectJob"); }
- 如果需要在算子内部获取密码(比如RichFlatMap、RichSinkFunction中),可以在算子的open方法中从运行时上下文读取全局参数:
@Override public void open(Configuration parameters) throws Exception { ParameterTool jobParams = (ParameterTool) getRuntimeContext() .getExecutionConfig() .getGlobalJobParameters(); String redisPassword = jobParams.getRequired("redis.password"); // 初始化Redis连接逻辑 }
更优的敏感参数传递方案
直接通过programArgs传明文密码存在一定泄露风险:参数值默认会展示在Flink Web UI的作业配置页、JobManager启动日志中,生产环境可以选择以下优化方案:
- 开启Flink的敏感配置掩码功能,将
redis.password加入敏感配置黑名单,UI和日志打印时会自动将对应值替换为******,避免明文泄露 - 如果企业内部有统一密钥管理服务,可以将本地解密后的临时密码写入密钥服务的临时存储路径,作业启动时通过内网拉取密码,避免密码随作业提交请求明文传输
- 有Kerberos等统一认证体系的场景,可以给Redis服务配置Kerberos认证,作业启动时自动通过票据认证,不需要手动传入密码
内容的提问来源于stack exchange,提问作者xdcsy
相关产品推荐
相关产品推荐

