使用Flink SQL进行双源表Join时如何读取RocksDB状态?
查询Flink SQL Join状态的方法及示例
你的场景是通过Flink SQL执行LEFT JOIN(基于id关联TABLE_1和TABLE_2),状态存储在RocksDB中,checkpoint目录指向/tmp。要查询这类Join产生的状态,推荐使用Flink State Processor API(官方支持的批处理式状态读取方案),以下是具体实现步骤和示例:
核心方案:State Processor API
该API允许以批处理方式读取、解析Flink作业(包括SQL作业)的状态数据,避免直接解析RocksDB底层文件(内部格式依赖强,不推荐)。
1. 准备依赖
确保项目引入对应Flink版本的State Processor相关依赖(以Flink 1.17.x为例):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-state-processor-api</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>1.17.1</version> </dependency>
2. 编写状态查询代码
SQL的LEFT JOIN会为左右表分别维护Keyed State(左表存储待匹配/已匹配数据,右表存储全量关联数据)。需先从Flink UI获取Join Operator的ID和状态名称,再编写代码读取:
import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.core.fs.Path; import org.apache.flink.state.api.SavepointReader; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.types.DataType; import org.apache.flink.types.Row; import org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend; public class QueryJoinState { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 替换为实际完成的checkpoint子目录(如/tmp/chk-1234,而非/tmp根目录) Path checkpointPath = new Path("/tmp/your-complete-checkpoint-dir"); SavepointReader savepointReader = Savepoint.load(env, checkpointPath, new EmbeddedRocksDBStateBackend(true)); // -------------------------- // 读取左表TABLE_1的Join状态 // -------------------------- DataType leftTableSchema = DataTypes.ROW( DataTypes.FIELD("headers", DataTypes.VARCHAR()), DataTypes.FIELD("id", DataTypes.VARCHAR()), DataTypes.FIELD("timestamp", DataTypes.TIMESTAMP_LTZ(3)), DataTypes.FIELD("type", DataTypes.VARCHAR()), DataTypes.FIELD("contentJson", DataTypes.VARCHAR()) ); TypeInformation<Row> leftRowType = (TypeInformation<Row>) leftTableSchema.getConversionClass(); savepointReader.readKeyedState( "JoinOperator", // 替换为Flink UI中查看的Join Operator ID "leftInputState", // 替换为UI中显示的左表状态名称 TypeInformation.of(String.class), // 关联Key的类型(对应id的VARCHAR) leftRowType ) .map(row -> String.format("左表ID: %s, 数据: %s", row.getField("id"), row)) .print("左表状态数据:"); // -------------------------- // 读取右表TABLE_2的Join状态 // -------------------------- DataType rightTableSchema = DataTypes.ROW( DataTypes.FIELD("headers", DataTypes.VARCHAR()), DataTypes.FIELD("id", DataTypes.VARCHAR()), DataTypes.FIELD("timestamp", DataTypes.TIMESTAMP_LTZ(3)), DataTypes.FIELD("type", DataTypes.VARCHAR()), DataTypes.FIELD("contentJson", DataTypes.VARCHAR()) ); TypeInformation<Row> rightRowType = (TypeInformation<Row>) rightTableSchema.getConversionClass(); savepointReader.readKeyedState( "JoinOperator", "rightInputState", // 替换为UI中显示的右表状态名称 TypeInformation.of(String.class), rightRowType ) .map(row -> String.format("右表ID: %s, 数据: %s", row.getField("id"), row)) .print("右表状态数据:"); env.execute("Query Flink SQL Join State"); } }
关键操作提示
- 获取Operator ID和状态名称:登录Flink Web UI,进入目标SQL作业的拓扑页面,找到Join算子,点击查看其状态详情,即可获取Operator ID和具体状态名称(如
leftInputState、rightInputState)。 - Checkpoint目录选择:必须使用已完成的checkpoint子目录(而非
/tmp根目录),根目录下的_metadata和shared是多checkpoint的共享元数据,无法直接读取。 - 类型匹配:代码中定义的Schema必须与SQL表字段类型严格一致,否则会出现解析失败。
快速查看状态统计(无需代码)
若仅需了解状态的条数、大小等统计信息,可直接在Flink Web UI中操作:
- 进入作业详情页,定位到Join算子
- 点击「State」标签,查看Keyed State的统计数据,包括状态名称、条目数、存储大小等
注意事项
- 禁止直接解析RocksDB底层文件:Flink对RocksDB的存储格式有内部实现依赖,版本更新可能导致格式变化,维护成本极高。
- 读取前确保作业状态稳定:优先读取已停止作业的checkpoint,避免读取运行中作业的checkpoint导致数据不一致。
内容的提问来源于stack exchange,提问作者hitesh
相关产品推荐
相关产品推荐

