You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Flink中外键关联的状态更新问题(DataStream/Table API)

背景

正在为所在组织进行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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 00:20:36