Apache Flink中外键关联的状态更新问题(DataStream/Table API)
Apache Flink 多流关联场景问题
背景
正在为所在组织进行Apache Flink多场景POC测试,目前处于学习阶段。我们当前使用Kafka Streams(搭配KTable)实现多流关联,但受限于其延迟问题,正在探索Apache Flink方案。
其中一个场景为外键关联:以员工(Employee)和部门(Department)为例,employee.deptId = department.deptId,一对多关系。
DataStream API 实现及问题
最初通过有状态流实现该关联,代码如下:
//DataStreamSource<Long> streamSource = env.fromSequence(1, 10000000); SingleOutputStreamOperator<Employee> employeeSourceStreamOperator = env .fromSource(employeeSource, WatermarkStrategy.noWatermarks(), "Employee") .map(value -> objectMapper.readValue(value, Employee.class)); SingleOutputStreamOperator<Department> departmentSingleOutputStreamOperator = env .fromSource(departmentSource, WatermarkStrategy.noWatermarks(), "Department") .map(value -> objectMapper.readValue(value, Department.class)); employeeSourceStreamOperator .connect(departmentSingleOutputStreamOperator) .keyBy(Employee::getDeptId, Department::getDeptId) .map(new RichCoMapFunction<Employee, Department, Object>() { // ListState to store multiple Employee instances for each department private ListState<Employee> employeeListState; // ValueState to store Department information private ValueState<Department> departmentState; @Override public void open(Configuration parameters) throws Exception { // Initialize ListState for Employee ListStateDescriptor<Employee> employeeListStateDescriptor = new ListStateDescriptor<>("employeeListState", Employee.class); employeeListState = getRuntimeContext().getListState(employeeListStateDescriptor); // Initialize ValueState for Department ValueStateDescriptor<Department> departmentStateDescriptor = new ValueStateDescriptor<>("departmentState", Department.class); departmentState = getRuntimeContext().getState(departmentStateDescriptor); } @Override public Object map1(Employee employee) throws Exception { // Process Employee stream // Store Employee information in MapState based on departmentId employeeListState.add(employee); // Try to join with Department information Department department = departmentState.value(); if (department != null) { // Join Employee and Department information return "Employee join:" + employee.getName() + " works in " + department.getDeptName(); } return ""; // No immediate result, need to wait for Department information } @Override public Object map2(Department department) throws Exception { // Process Department stream // Store Department information in ValueState departmentState.update(department); // Try to join with Employee information Iterable<Employee> employees = employeeListState.get(); if (employees != null) { // Join Employee and Department information for each employee in the list StringBuilder result = new StringBuilder(); for (Employee employee : employees) { result.append("Department join:" + generateOutput(employee, department)).append("\n"); } return result.toString(); } return ""; // No immediate result, need to wait for Employee information } private String generateOutput(Employee employee, Department department) { return employee.getName() + " works in " + department.getDeptName(); } }) .sinkTo(new PrintSink<>());
该实现可正常运行,但遇到员工换部门的场景时,新关联易实现,但难以从原部门状态中移除该员工。
Table API 实现及测试问题
后续尝试用Table API实现,代码如下:
SingleOutputStreamOperator<Employee> employeeSourceStreamOperator = env .fromSource(employeeSource, WatermarkStrategy.noWatermarks(), "Employee") .map(value -> objectMapper.readValue(value, Employee.class)); Table employeeTable = tableEnvironment .fromDataStream(employeeSourceStreamOperator, $("deptId").as("empDeptId"),$("name").as("name")); SingleOutputStreamOperator<Department> departmentSingleOutputStreamOperator = env .fromSource(departmentSource, WatermarkStrategy.noWatermarks(), "Department") .map(value -> objectMapper.readValue(value, Department.class)); Table departmentTable = tableEnvironment .fromDataStream(departmentSingleOutputStreamOperator, Schema.newBuilder().build()); Table table = employeeTable.join(departmentTable) .where($("empDeptId").isEqual($("deptId"))) .select($("name"), $("deptName")); DataStream<Row> rowDataStream = tableEnvironment.toChangelogStream(table); rowDataStream.sinkTo(new PrintSink<>());
测试流程及异常结果
- 发布部门1及员工6数据,输出
5> +I[Employee 6, Department 1] - 发布员工6换部门至2的更新,输出
5> +I[Employee 6, Department 2] - 发布部门1的更新,仍输出
5> +I[Employee 6, Department 1],员工6仍关联原部门。
测试仓库包含测试数据生成工具EmployeeDepartmentTestDataGenerator.java。
请问我遗漏了什么配置或逻辑?
内容的提问来源于stack exchange,提问作者Umesh Kamat
相关产品推荐
相关产品推荐

