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

使用Flink SQL进行双源表Join时如何读取RocksDB状态?

你的场景是通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 15:01:17