Flink中Filter算子里jobStartTimeStamp恒为0的原因及解决方法
问题分析与解决:Flink静态全局变量在算子中值为0的问题
原因分析
Flink是分布式计算框架,作业执行存在跨进程隔离:
main函数运行在提交作业的客户端JVM进程中,你在这里给静态变量jobStartTimeStamp赋值,这个值仅在客户端进程内有效。- 算子逻辑(比如代码里的
FilterFunction匿名内部类)会被序列化后发送到TaskManager的独立工作进程执行。TaskManager进程重新加载你的类时,静态变量会使用Java默认初始值0,客户端的赋值完全不会同步到TaskManager进程中。
这就是Filter算子中jobStartTimeStamp始终为0的核心原因——静态变量的作用域仅限单个JVM进程,跨进程的分布式场景下无法共享客户端的赋值结果。
解决办法
方法1:使用局部final变量传递(最简单)
将jobStartTimeStamp改为main函数内的局部变量并声明为final(或Java 8+的"有效final"变量,即后续不再修改),匿名内部类会捕获该变量的值并序列化携带到TaskManager进程:
public static void main(String[] args) { // 替换静态全局变量为局部final变量 final long jobStartTimeStamp = System.currentTimeMillis() + 120000; System.out.println("StartTime" + jobStartTimeStamp); ... if (isOutputTagEnable) { SinkFunction<MonitorDataHolder> lateMonitorDataProducer = new FlinkKafkaProducer010<>(properties.getProperty("lateMonitorDataSinkTopic", "HaMonitor_MonitorDataDiscard"), new LateMonitorDataSerializationSchema(), kafkaSinkConfig); monitorDataCompute .getSideOutput(outputTag) .filter(new FilterFunction<MonitorDataHolder>() { @Override public boolean filter(MonitorDataHolder value) throws Exception { // 直接使用捕获的final变量 return value.getCurrentTimeStamp() > jobStartTimeStamp; } }) .addSink(lateMonitorDataProducer) .setParallelism(Integer.parseInt(properties.getProperty("lateMonitorDataSinkParallelism", "3"))) .uid("SinkLateMonitorData"); } }
方法2:通过Flink Configuration传递
将时间值存入Flink配置,算子通过RuntimeContext在初始化阶段读取配置值:
public static void main(String[] args) { long jobStartTimeStamp = System.currentTimeMillis() + 120000; System.out.println("StartTime" + jobStartTimeStamp); // 将时间值写入Configuration Configuration config = new Configuration(); config.setLong("job.start.time", jobStartTimeStamp); // 创建执行环境时传入配置 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config); ... if (isOutputTagEnable) { SinkFunction<MonitorDataHolder> lateMonitorDataProducer = new FlinkKafkaProducer010<>(properties.getProperty("lateMonitorDataSinkTopic", "HaMonitor_MonitorDataDiscard"), new LateMonitorDataSerializationSchema(), kafkaSinkConfig); monitorDataCompute .getSideOutput(outputTag) .filter(new FilterFunction<MonitorDataHolder>() { private transient long jobStartTime; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 从配置中读取时间值 jobStartTime = parameters.getLong("job.start.time", 0L); } @Override public boolean filter(MonitorDataHolder value) throws Exception { return value.getCurrentTimeStamp() > jobStartTime; } }) .addSink(lateMonitorDataProducer) .setParallelism(Integer.parseInt(properties.getProperty("lateMonitorDataSinkParallelism", "3"))) .uid("SinkLateMonitorData"); } }
方法3:使用广播变量(多算子共享场景)
如果需要给多个算子传递该时间值,可以用广播变量实现跨TaskManager节点的变量共享:
public static void main(String[] args) { long jobStartTimeStamp = System.currentTimeMillis() + 120000; System.out.println("StartTime" + jobStartTimeStamp); // 将时间值包装为广播数据集 List<Long> broadcastData = Collections.singletonList(jobStartTimeStamp); ... if (isOutputTagEnable) { SinkFunction<MonitorDataHolder> lateMonitorDataProducer = new FlinkKafkaProducer010<>(properties.getProperty("lateMonitorDataSinkTopic", "HaMonitor_MonitorDataDiscard"), new LateMonitorDataSerializationSchema(), kafkaSinkConfig); monitorDataCompute .getSideOutput(outputTag) // 注册广播变量 .broadcast(broadcastData) .filter(new FilterFunction<MonitorDataHolder>() { private transient long jobStartTime; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 获取广播变量的值 jobStartTime = getRuntimeContext().getBroadcastVariable(0).get(0); } @Override public boolean filter(MonitorDataHolder value) throws Exception { return value.getCurrentTimeStamp() > jobStartTime; } }) .addSink(lateMonitorDataProducer) .setParallelism(Integer.parseInt(properties.getProperty("lateMonitorDataSinkParallelism", "3"))) .uid("SinkLateMonitorData"); } }
内容的提问来源于stack exchange,提问作者brotherxu
相关产品推荐
相关产品推荐

