Apache Flink批处理模式下基于DataSet约束的条件处理实现
问题背景
在Flink批处理场景中,读取CSV文件转为DataSet<MyObject>后完成数据校验,每条记录的isValid字段会被标记为true(有效)或false(无效)。需要判断所有记录是否全部有效,并据此分支处理:全有效则执行后续转换并生成Ack文件,否则生成Nck文件。但存在以下限制:
- 不能使用
collect()/count()等急切执行函数(detached模式下会触发报错) - 无法在作业运行前直接获取DataSet中的数据状态
- 不能通过常规方式从DataSet中直接提取
MyObject的状态
触发的报错信息:
Caused by: org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: Job was submitted in detached mode. Results of job execution, such as accumulators, runtime, etc. are not available. Please make sure your program doesn't call an eager execution function [collect, print, printToErr, count].
MyObject定义:
public class MyObject{ private FileDetails fileDetails; private String fileName; public boolean isValid ; public String invalidReason; }
原示例代码:
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); String fileName = "myFile.csv"; DataSet<MyObject> readCsvData = readAndProcessCsvData(env, fileName); DataSet<MyObject> validateFile = validateFile(readCsvData); if(toCheckIsValidFlag){ //********** process further *********** generateAckFile(); } else { generateNckFile(); } env.execute();
解决方案:使用Flink累加器(Accumulator)+ 作业执行后获取结果
Flink的累加器可以在作业运行过程中统计全局状态,作业执行完成后可获取累加器结果,完全规避急切函数的使用。具体实现步骤如下:
1. 定义自定义累加器
实现一个累加器来统计无效记录的总数:
public class InvalidRecordCounter implements Accumulator<Boolean, Integer> { private int count = 0; @Override public void add(Boolean isValid) { if (!isValid) { count++; } } @Override public Integer getLocalValue() { return count; } @Override public void resetLocal() { count = 0; } @Override public void merge(Accumulator<Boolean, Integer> other) { this.count += other.getLocalValue(); } @Override public Accumulator<Boolean, Integer> clone() { InvalidRecordCounter counter = new InvalidRecordCounter(); counter.count = this.count; return counter; } }
2. 在数据校验环节注册并使用累加器
在validateFile函数中,遍历记录时同步更新累加器:
public static DataSet<MyObject> validateFile(DataSet<MyObject> inputData, ExecutionEnvironment env) { InvalidRecordCounter invalidCounter = new InvalidRecordCounter(); // 注册累加器,指定名称方便后续获取结果 env.registerAccumulator("invalidRecordCount", invalidCounter); return inputData.map(new MapFunction<MyObject, MyObject>() { @Override public MyObject map(MyObject value) throws Exception { // 替换为实际业务校验逻辑,更新isValid和invalidReason if (value.getFileName() == null || value.getFileName().isEmpty()) { value.isValid = false; value.invalidReason = "File name is empty"; } else { value.isValid = true; } // 将当前记录的有效性状态传入累加器 invalidCounter.add(value.isValid); return value; } }); }
3. 执行作业后获取累加器结果并分支处理
提交作业后,通过JobExecutionResult获取累加器统计值,再判断分支逻辑:
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); String fileName = "myFile.csv"; DataSet<MyObject> readCsvData = readAndProcessCsvData(env, fileName); DataSet<MyObject> validateFile = validateFile(readCsvData, env); // 先将校验后的数据写入临时存储(按需调整) validateFile.writeAsCsv("path/to/temp/validated-records"); // 执行作业并获取执行结果 JobExecutionResult result = env.execute("CSV Data Validation Job"); // 从结果中提取无效记录总数 int invalidCount = result.getAccumulatorResult("invalidRecordCount"); if (invalidCount == 0) { // 所有记录有效,执行后续处理并生成Ack文件 processFurther(validateFile); generateAckFile(); } else { // 存在无效记录,生成Nck文件 generateNckFile(); }
补充优化:分离有效/无效数据输出
如果需要将有效、无效数据分别输出到不同文件,可以在校验后添加filter操作:
// 分离有效和无效数据集 DataSet<MyObject> validData = validateFile.filter(obj -> obj.isValid); DataSet<MyObject> invalidData = validateFile.filter(obj -> !obj.isValid); // 分别下沉到对应文件 validData.writeAsCsv("path/to/output/valid-data"); invalidData.writeAsCsv("path/to/output/invalid-data");
内容的提问来源于stack exchange,提问作者Dhanaji Patil

